//! TRACES: FR-EXP-6 | FR-EXP-7 | FR-NC-10 //! Placing an exported file, and the cache that makes offline unremarkable. //! //! [`dr_export`] turns a frame into bytes and a name and stops there, because //! where those bytes go differs by more than a path. This is the other half: //! it decides the destination and gets them there. //! //! # Everything is staged first //! //! An export to the server is written to a local **outbox** before any upload //! is attempted, and the upload drains that outbox afterwards. Not a fallback //! for the offline case — the *only* path, with offline merely meaning the //! drain finds nothing to do. //! //! Doing it the other way, uploading directly and staging only on failure, //! looks simpler and has two bad properties. The failure path is then the one //! that is rarely exercised and always broken, and the moment the network //! drops mid-batch some exports exist and some do not with nothing recording //! which. Staging first means an export is *finished* the instant it is //! written; the upload is a separate promise the app keeps later. //! //! # Why the outbox is not in the cache //! //! It lives beside the catalog, with the thumbnail shards, rather than under //! the evictable cache. `dr_catalog::cache` draws the line already: passively //! cached originals are a convenience and go under LRU, pinned ones are a //! promise and never do. An export waiting to upload is a promise — the user //! was told the export succeeded — and sweeping it away to reclaim disk would //! destroy work that no longer exists anywhere else. //! //! # The batch (FR-EXP-7) //! //! The second half of this file is a worker that exports a whole selection. //! It runs on a thread of its own and reaches the interface through the same //! two mechanisms everything else here does: an `mpsc` channel drained by a //! Slint timer, and a row in [`crate::activity`]. //! //! **Why the worker opens its own sessions.** A [`crate::DevelopSession`] owns a GPU //! `AdjustPass`, and the one the interface is holding is the one the canvas //! renders from — handing it to a worker would mean the develop view could not //! draw while a batch ran, which is the freeze this exists to remove. A //! `GpuContext` is an `Arc` pair over a device and queue and is cheap to //! clone, so the worker takes a clone and opens each photograph for itself. //! The cost is a `Demosaicer` and an `AdjustPass` per image rather than one //! for the run; against a full-resolution decode, render and encode it is //! small, and it keeps this file out of the pipeline that `develop` owns. //! //! **The open image is the exception.** Its edit lives in the interface's //! session and may not have reached a sidecar yet, so a worker that re-opened //! the file for itself would export the *saved* version rather than the one on //! screen. That one frame is therefore rendered by the caller and handed over //! as [`Source::Rendered`]; everything after the render still moves off the UI //! thread. It travels with the header its session was opened from, so an //! export made from the develop button discloses exactly what an export of the //! same file from the grid does — subject, both times, to what the settings //! allow (FR-EXP-8). use crate::executors::{self, Executor}; use std::collections::{HashMap, HashSet}; use std::path::{Path, PathBuf}; use std::rc::Rc; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender}; use std::sync::Arc; use std::time::Duration; use dr_export::{Encoded, NameContext}; use dr_sync::RemotePath; use dr_sync::{Account, Connection}; use dr_types::{CollisionPolicy, ExportSettings, ExportTarget}; use slint::ComponentHandle as _; use crate::{AppWindow, Library}; /// Where exports wait for a server that is not there yet. /// /// Beside the catalog, for the reason in the module docs. Per account, /// because the destination folder is a path on one particular server and an /// entry queued for one account is meaningless to another. pub fn outbox_dir(account: &Account) -> PathBuf { crate::library::catalog_path(account) .parent() .map(|p| p.join("outbox")) .unwrap_or_else(|| std::env::temp_dir().join("darkroom-outbox")) } /// One export waiting to go up. /// /// The record sits beside the bytes as `.dest`, holding the remote /// folder it belongs in. A single flat file rather than a database: the queue /// is small, the entries are independent, and the recovery story for a /// half-written text file is to ignore it — which is exactly what parsing it /// does. #[derive(Debug, Clone, PartialEq, Eq)] pub struct Pending { /// The staged bytes on this device. pub local: PathBuf, /// Remote folder, relative to the library root. Empty means the root. pub remote_dir: String, /// TRACES: FR-EXP-10 /// `remote_dir` is relative to the *account* root instead: an album on /// the server lives outside the library, where a scan would not /// catalogue its JPEGs as photographs. A third line in the record, so a /// record written before albums — two lines, a stray `/` and all — keeps /// the meaning it was written with. pub account: bool, /// The filename to give it there. pub name: String, /// What to do if the server already holds `name`. Decided when the /// export was made, applied when it uploads: the batch could not see /// the server, so the name it chose is a request, not a fact. `None` /// for a record written before the policy travelled with it, and for /// a merge's composite, which the catalog has already recorded under /// its name — both are sent as named. pub collision: Option, } impl Pending { /// Full remote path for this entry, under `root`. fn remote_path(&self, root: &str) -> RemotePath { let mut parts: Vec<&str> = Vec::new(); for segment in [self.base(root), self.remote_dir.as_str()] { for part in segment.split('/') { if !part.is_empty() { parts.push(part); } } } parts.push(&self.name); RemotePath::new(parts.join("/")) } /// The folder this entry's file belongs in, as a remote path. fn remote_folder(&self, root: &str) -> RemotePath { let mut parts: Vec<&str> = Vec::new(); for segment in [self.base(root), self.remote_dir.as_str()] { for part in segment.split('/') { if !part.is_empty() { parts.push(part); } } } RemotePath::new(parts.join("/")) } /// What `remote_dir` is relative to: the library root, or the account's. fn base<'a>(&self, root: &'a str) -> &'a str { if self.account { "" } else { root } } } /// 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(), collision: None, } .remote_path(root) } /// Where an export was put, for the interface to report. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Placed { /// Written straight to a folder on this device. Device(PathBuf), /// Staged locally, awaiting upload to the named remote folder as `name`. /// The staged copy may be called something else — see [`stage`]. Queued { local: PathBuf, remote_dir: String, name: String, }, } impl Placed { /// The file's name where it was asked to go — what an album records. pub fn file_name(&self) -> String { match self { Placed::Device(path) => path .file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or_default(), Placed::Queued { name, .. } => name.clone(), } } /// A sentence for the status line. pub fn describe(&self) -> String { match self { Placed::Device(path) => format!( "Exported {}", path.file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or_default() ), // Named as queued rather than exported: the file is real and // finished, but it is not yet where the user asked for it, and // saying "exported to the library" before it has uploaded would be // a claim the app cannot keep if the disk is pulled. Placed::Queued { remote_dir, .. } => { let dir = if remote_dir.is_empty() { "the library root".to_string() } else { remote_dir.trim_start_matches('/').to_string() }; format!("Queued for {dir}") } } } } /// Write an encoded export to wherever the settings say it goes. /// /// `destination` is a filesystem path for [`ExportTarget::Device`] and a /// remote folder for [`ExportTarget::Remote`] — the widening `ExportSettings` /// documents, resolved here because this is the layer that knows what a path /// means on this platform. pub fn place( encoded: &Encoded, target: ExportTarget, destination: &str, outbox: &Path, collision: CollisionPolicy, ) -> Result { match target { ExportTarget::Device => { if destination.trim().is_empty() { return Err("No export folder is set. Choose an album to export to.".into()); } // TRACES: FR-EXP-10 // A SAF tree on Android: written through the provider, which may // rename on a collision, so the name it reports is the one kept. #[cfg(target_os = "android")] if destination.starts_with("content://") { let replace = collision_replaces(&encoded.name, destination); let written = crate::saf::write( destination, &encoded.name, mime_for(&encoded.name), &encoded.bytes, replace, )?; return Ok(Placed::Device(PathBuf::from(written))); } let dir = PathBuf::from(destination); std::fs::create_dir_all(&dir).map_err(|e| format!("{}: {e}", dir.display()))?; let path = dir.join(&encoded.name); std::fs::write(&path, &encoded.bytes) .map_err(|e| format!("{}: {e}", path.display()))?; Ok(Placed::Device(path)) } ExportTarget::Remote => { let local = stage(encoded, destination, outbox, collision)?; Ok(Placed::Queued { local, remote_dir: destination.to_string(), name: encoded.name.clone(), }) } } } /// Whether writing `name` into a SAF folder should replace a file already /// there. By the time `place` runs, the batch has already chosen the name /// under the collision policy, so a name that is taken can only have been /// chosen under Overwrite. #[cfg(target_os = "android")] fn collision_replaces(name: &str, tree: &str) -> bool { crate::saf::exists(tree, name) } /// The MIME type a document provider is told, from the name the encoder gave. #[cfg(target_os = "android")] fn mime_for(name: &str) -> &'static str { match Path::new(name) .extension() .and_then(|e| e.to_str()) .map(str::to_ascii_lowercase) .as_deref() { Some("png") => "image/png", Some("tif" | "tiff") => "image/tiff", _ => "image/jpeg", } } /// Write bytes and their destination record into the outbox. fn stage( encoded: &Encoded, remote_dir: &str, outbox: &Path, collision: CollisionPolicy, ) -> Result { std::fs::create_dir_all(outbox).map_err(|e| format!("{}: {e}", outbox.display()))?; // The staged name is the export's name, deduplicated against the outbox // rather than against the server: two exports queued before either has // uploaded would otherwise overwrite each other here, and the second one // would silently replace the first before anybody saw it. let mut candidate = outbox.join(&encoded.name); let mut n = 1; while candidate.exists() { let stem = Path::new(&encoded.name) .file_stem() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_else(|| "export".into()); let ext = Path::new(&encoded.name) .extension() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_default(); candidate = outbox.join(format!("{stem}-{n}.{ext}")); n += 1; if n > 10_000 { return Err("the outbox is full of files by this name".into()); } } // Bytes first, then the record. The order matters on a process that may // be killed between the two: an orphan payload with no record is ignored // by the drain and swept later, where a record naming bytes that were // never written would be a permanent failure retried forever. std::fs::write(&candidate, &encoded.bytes) .map_err(|e| format!("{}: {e}", candidate.display()))?; let record = destination_record(&candidate); // The remote folder and the intended name, one per line. Not JSON: two // strings do not need a parser, and a format a human can repair by hand // is worth something for a queue holding the only copy of someone's work. // A folder spelled from `/` is an album's, relative to the account; it // is recorded as such on a third line — see `Pending::account`. The // fourth is the collision policy, which an older build stops before // reading and so uploads as named, as it always did. let (dir, base) = match remote_dir.strip_prefix('/') { Some(dir) => (dir, "account"), None => (remote_dir, "library"), }; let text = format!( "{dir}\n{}\n{base}\n{}\n", encoded.name, policy_word(collision) ); std::fs::write(&record, text).map_err(|e| format!("{}: {e}", record.display()))?; Ok(candidate) } /// The record that names where a staged payload goes: `name.ext.dest` /// beside it. One function so the merge, which stages its own composite, /// and the drain agree on the name. pub(crate) fn destination_record(payload: &Path) -> PathBuf { payload.with_extension(format!( "{}.dest", payload .extension() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_default() )) } /// 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 /// order in [`stage`]. pub fn pending(outbox: &Path) -> Vec { let Ok(entries) = std::fs::read_dir(outbox) else { return Vec::new(); }; let mut out = Vec::new(); for entry in entries.flatten() { let path = entry.path(); if path.extension().and_then(|e| e.to_str()) != Some("dest") { continue; } // `photo.jpg.dest` describes `photo.jpg`. let local = path.with_extension(""); if !local.exists() { continue; } let Ok(text) = std::fs::read_to_string(&path) else { continue; }; let mut lines = text.lines(); let remote_dir = lines.next().unwrap_or("").to_string(); let name = lines.next().unwrap_or("").to_string(); let account = lines.next() == Some("account"); let collision = lines.next().and_then(policy_from_word); if name.is_empty() { continue; } out.push(Pending { local, remote_dir, name, account, collision, }); } // Stable order so a drain is reproducible and a stuck entry is obvious // rather than appearing to move around the queue. out.sort_by(|a, b| a.local.cmp(&b.local)); out } /// A collision policy as an outbox record spells it. fn policy_word(policy: CollisionPolicy) -> &'static str { match policy { CollisionPolicy::Overwrite => "overwrite", CollisionPolicy::Skip => "skip", CollisionPolicy::Increment => "increment", } } fn policy_from_word(word: &str) -> Option { match word { "overwrite" => Some(CollisionPolicy::Overwrite), "skip" => Some(CollisionPolicy::Skip), "increment" => Some(CollisionPolicy::Increment), _ => None, } } /// How many exports are waiting. For the interface to show, and cheap enough /// to call on a redraw. pub fn pending_count(outbox: &Path) -> usize { pending(outbox).len() } /// 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. #[derive(Debug)] 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, /// The server held its name and the policy was Skip. Skipped, } /// What each destination folder holds, as far as this drain knows: listed /// once on the first entry bound for it, and added to as files land. One /// request per folder per drain rather than one per file — a batch of /// three hundred into one album is one listing. type Listings = HashMap>; /// 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, listings: &mut Listings, ) -> 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. let folder = entry.remote_folder(root); backend .create_dir(&folder) .await .map_err(|e| e.to_string())?; // TRACES: FR-EXP-6 // The batch named this file without seeing the server, so the policy it // was made under is applied here, against what the folder holds now. // Overwrite, and a record that carries no policy, put it as named. let names = match entry.collision { Some(CollisionPolicy::Increment | CollisionPolicy::Skip) => { let key = folder.as_str().to_string(); if !listings.contains_key(&key) { let listed = backend .list(&folder, None) .await .map_err(|e| e.to_string())?; let held = listed.iter().map(|e| e.path.name().to_string()).collect(); listings.insert(key.clone(), held); } listings.get_mut(&key) } _ => None, }; let mut renamed = None; if let Some(held) = names.as_deref() { if held.contains(&entry.name) { if entry.collision == Some(CollisionPolicy::Skip) { log::info!("upload: {} is already on the server; skipped", entry.name); clear(entry); return Ok(Sent::Skipped); } let free = step_past(&entry.name, &|n| held.contains(n)) .ok_or_else(|| format!("{}: ten thousand names taken", entry.name))?; renamed = Some(Pending { name: free, ..entry.clone() }); } } let sent = renamed.as_ref().unwrap_or(entry); backend .put(&sent.remote_path(root), bytes, None) .await .map_err(|e| e.to_string())?; if let Some(held) = names { held.insert(sent.name.clone()); } if renamed.is_some() { log::info!( "upload: {} was taken on the server; sent as {}", entry.name, sent.name ); rename_in_album(library, entry, &sent.name); } let thumbnails = take_thumbnails(&entry.local); register_upload(backend, library, root, sent, thumbnails).await; clear(entry); Ok(Sent::Uploaded) } /// TRACES: FR-EXP-10 /// Tell the album a file it recorded arrived under another name. Only an /// album's entries are recorded by name; a library export is the scan's. fn rename_in_album(library: &LibraryFiles, entry: &Pending, to: &str) { if !entry.account { return; } let result = dr_catalog::Catalog::open(&library.catalog).and_then(|catalog| { dr_catalog::albums::rename_export(catalog.connection(), &entry.remote_dir, &entry.name, to) }); if let Err(e) = result { log::warn!("upload: recording {} as {to} in its album: {e}", entry.name); } } /// 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(()); /// Drain the outbox to the server. /// /// Its own thread with its own runtime, like every other network path here — /// the Slint loop must never block (NFR-P9). /// /// A failure leaves the entry in place and stops the run. Continuing past a /// network error would burn the whole queue against a server that is not /// answering, and the next pass costs nothing. pub fn spawn_upload( conn: Connection, root: String, outbox: PathBuf, ) -> std::sync::mpsc::Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "upload", move || { // One drain at a time. Three things start one — a sync pass, an // export, a finished merge — and they used to run side by side over // the same queue: each read the same 800 MB composite into memory and // sent it, and the one that came second found its record cleared // under it and reported the file missing. Waiting here costs nothing: // the drain that holds the lock is already sending everything queued, // and the one that waited finds the queue empty and finishes. let _draining = OUTBOX_DRAIN .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); let rt = match crate::net_runtime::build() { Ok(rt) => rt, Err(e) => { let _ = tx.send(UploadMessage::Finished { uploaded: 0, remaining: pending_count(&outbox), error: Some(e.to_string()), landed: 0, }); return; } }; rt.block_on(async { let backend = match crate::remote::connect(&conn) { Ok(b) => b, Err(e) => { let _ = tx.send(UploadMessage::Finished { 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; let mut listings = Listings::new(); for (i, entry) in queue.iter().enumerate() { let _ = tx.send(UploadMessage::Status(format!( "uploading {} ({}/{total})", entry.name, i + 1 ))); match send(&*backend, &library, &root, entry, &mut listings).await { Ok(Sent::Gone | Sent::Skipped) => {} Ok(Sent::Uploaded) => { uploaded += 1; if !entry.account { landed += 1; } } Err(e) => { error = Some(e); break; } } } let _ = tx.send(UploadMessage::Finished { uploaded, remaining: pending_count(&outbox), error, landed, }); }); }); rx } // --------------------------------------------------------------------------- // Batch export (FR-EXP-7) // --------------------------------------------------------------------------- /// TRACES: NFR-ARCH-3 /// The stop flag a running batch reads. /// /// One bit rather than a channel. Cancellation is read thousands of times more /// often than it is written, and a message would only be seen if the worker /// happened to be listening at the moment the button was pressed rather than at /// the next point it is safe to stop. /// /// `Relaxed` throughout: the flag guards no other data, so there is nothing for /// an acquire/release pair to publish. The only ordering that matters is that /// the write eventually becomes visible, which every ordering guarantees. #[derive(Clone, Default)] pub struct Cancel(Arc); impl Cancel { pub fn cancel(&self) { self.0.store(true, Ordering::Relaxed); } pub fn is_cancelled(&self) -> bool { self.0.load(Ordering::Relaxed) } } /// How long the worker blocks on another worker before rereading the flag. /// /// TRACES: NFR-ARCH-3 /// This is the cancellation bound for everything the batch spends its time /// *waiting* on — a download of tens of megabytes would otherwise hold a /// cancelled batch open until the transfer finished. It is not the bound for /// the frame being rendered and encoded when the button is pressed: that has no /// interior stopping point, so the true worst case is one image, and on a 24 MP /// frame that is well past NFR-ARCH-3's 100 ms target. Closing that gap needs /// the render itself to become interruptible (NFR-ARCH-2's scheduler), not a /// finer poll here. const CANCEL_POLL: Duration = Duration::from_millis(100); /// How often the drain looks at what the worker has said. const DRAIN_INTERVAL: Duration = Duration::from_millis(120); /// Where one image's pixels come from. pub enum Source { /// A photograph in the library, fetched and rendered by the worker. /// /// Carries its own cache context because that is assembled from the catalog /// and the session, and a worker thread can reach neither. Library { path: String, cache: Option, }, /// A frame the caller has already rendered — the image open in develop. /// /// See the module docs for why this one cannot be left to the worker. Rendered { stem: String, /// TRACES: FR-EXP-8 /// What the photograph's own file said about itself, as the session /// remembers it — see [`crate::DevelopSession::source_metadata`]. /// /// The header travels rather than a finished /// [`dr_export::SourceMetadata`], so the allowlist below stays the one /// place a source field becomes exportable, and the two arms of /// [`export_one`] reach the encoder through the same transcription. /// /// `None` where the file had no header to read, which is an ordinary /// outcome and not an error. header: Option, frame: dr_export::Frame, }, } impl Source { /// What to call this image in a progress row. fn describe(&self) -> String { match self { Source::Library { path, .. } => path.rsplit('/').next().unwrap_or(path).to_string(), Source::Rendered { stem, .. } => stem.clone(), } } } /// Everything a batch needs, assembled on the UI thread. /// /// Assembled there and not in the worker because most of it is only reachable /// from a `RefCell` the interface owns: the session, the settings record, the /// catalog. Passing the finished answers means the worker borrows nothing. pub struct BatchRequest { pub sources: Vec, /// TRACES: FR-EXP-10 /// The catalog image behind each source, by position, where there is one /// — what an album records its files against. Shorter than `sources`, or /// empty, where the caller does not know. pub images: Vec>, /// Credentials for the account the library is open on. `None` where no /// library is open, which is fine for a [`Source::Rendered`] and fatal for /// anything that has to be fetched. pub conn: Option, pub settings: ExportSettings, /// TRACES: FR-EXP-6 /// Names a server destination is known to hold — what the album records /// of earlier exports. A queued export cannot ask the server, so this is /// what [`resolve_batch_name`] checks instead; the upload checks the /// server itself for anything put there some other way. Unused for a /// device export, which looks at the folder. pub remote_names: HashSet, pub outbox: PathBuf, pub sidecar_cache: PathBuf, pub offline: bool, /// `None` on a build with no adapter, where a library image cannot be /// rendered at all — reported per image rather than refused up front, so /// the reason lands in the same place every other failure does. pub gpu: Option, /// TRACES: FR-RAW-2 /// What reads the header and decodes the sensor data, chosen by whoever /// assembled the batch; the worker never names one. pub decoder: &'static dyn dr_decode::Decoder, } /// TRACES: NFR-ARCH-4 /// Why one image of a batch produced no file. /// /// One variant per stage. "Export failed" repeated forty times is not a report /// anybody can act on: a destination folder that cannot be written is one fix, /// a body rawler cannot decode is another, and a server that went away is a /// third — and only the first is worth stopping the batch to correct. #[derive(Debug, thiserror::Error)] pub enum ItemError { #[error("could not be fetched: {0}")] Fetch(String), #[error("could not be opened for editing: {0}")] Open(String), #[error("could not be rendered: {0}")] Render(String), /// The name was taken and the collision policy is Skip. Not really a /// failure — the user asked for exactly this — but it is still an image /// that produced no file, and a batch reporting "40 exported" when it wrote /// 31 would be lying. #[error("a file of that name is already there, and the collision setting is Skip")] NameTaken, #[error(transparent)] Encode(#[from] dr_export::ExportError), #[error("could not be written: {0}")] Place(String), } /// What the worker says as it goes. #[derive(Debug)] pub enum BatchMessage { /// Work has begun on one image; `done` counts the ones behind it. Started { done: usize, total: usize, name: String, }, /// One image is finished, for better or worse. Item { name: String, /// The catalog image it came from, where the request said. image: Option, outcome: Result, }, Finished { exported: usize, failed: usize, cancelled: bool, }, } /// Export a selection on a thread of its own. pub fn spawn_batch(request: BatchRequest, cancel: Cancel) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::GpuSubmit, "export", move || { run(request, &cancel, &tx) }); rx } /// The batch itself, one image at a time. /// /// Sequential rather than parallel, which FR-EXP-7's "uses all available cores" /// does not yet get. One reason and one excuse: the reason is that a /// full-resolution frame is tens of megabytes and four in flight is a /// straightforward way to exhaust a tablet; the excuse is that the GPU is /// shared with the interface, and the render is where the time goes. fn run(mut request: BatchRequest, cancel: &Cancel, tx: &Sender) { let sources = std::mem::take(&mut request.sources); let total = sources.len(); // Names this run has already written. See [`resolve_batch_name`] for why // the destination alone is not enough to keep two exports apart. let mut issued: HashSet = HashSet::new(); let (mut exported, mut failed) = (0usize, 0usize); // What earlier batches queued for the same folder and has not uploaded // yet is as much in the way as what is already there. if request.settings.target == ExportTarget::Remote { let names = queued_for(&request.outbox, &request.settings.destination); request.remote_names.extend(names); } for (i, source) in sources.into_iter().enumerate() { if cancel.is_cancelled() { break; } let name = source.describe(); let _ = tx.send(BatchMessage::Started { done: i, total, name: name.clone(), }); // `None` is cancellation mid-image, which is not an outcome for this // photograph: nothing failed, the user simply stopped asking. let Some(outcome) = export_one(&request, source, i as u32 + 1, &mut issued, cancel) else { break; }; // Counted, reported, and then on to the next one. A failure here must // not end the run — the whole point of exporting three hundred frames // unattended is that the one unreadable file is a line in a report // rather than an evening lost (FR-EXP-7). match &outcome { Ok(_) => exported += 1, Err(_) => failed += 1, } let image = request.images.get(i).copied().flatten(); let _ = tx.send(BatchMessage::Item { name, image, outcome, }); } let _ = tx.send(BatchMessage::Finished { exported, failed, cancelled: cancel.is_cancelled(), }); } /// The names waiting in the outbox for `destination`, spelled as a batch /// spells a remote folder: from `/` for an album, relative to the library /// otherwise. fn queued_for(outbox: &Path, destination: &str) -> impl Iterator { let account = destination.starts_with('/'); let folder = destination.trim_matches('/').to_string(); pending(outbox) .into_iter() .filter(move |p| p.account == account && p.remote_dir.trim_matches('/') == folder) .map(|p| p.name) } /// One image, start to finish. `None` where the run was cancelled part way. fn export_one( request: &BatchRequest, source: Source, sequence: u32, issued: &mut HashSet, cancel: &Cancel, ) -> Option> { // TRACES: FR-EXP-8 // The second and fourth elements are what the photograph's own file said // about itself: the capture date the `{date}` token names, and the record // the encoder may write from. A library image is decoded here and a // rendered one arrives with the header its session was opened from, and // both go through [`from_header`] — the same reading, so the same // photograph exported from the two buttons cannot come out saying // different things about itself. let (stem, date, frame, source_metadata) = match source { Source::Rendered { stem, header, frame, } => { // A file with no header stays a file with no header: an empty // `{date}` and nothing for the encoder to copy. Substituting // today's date here would put a lie in the filename of every // scanned negative. let (date, carried) = match header { Some(meta) => { let (date, carried) = from_header(&meta); (date, Some(carried)) } None => (String::new(), None), }; (stem, date, frame, carried) } Source::Library { path, cache } => { match render_from_library(request, &path, cache, cancel)? { Ok(rendered) => rendered, Err(e) => return Some(Err(e)), } } }; Some(place_frame( request, &stem, &date, sequence, &frame, source_metadata.as_ref(), issued, )) } /// TRACES: FR-EXP-8 /// What a source file's header contributes to an export. /// /// Two things, and they are read together because they come from the same /// place: the `{date}` token, and the record the encoder writes from. One /// function for both kinds of source, so that a develop export and a library /// export of the same photograph cannot disagree about the day it was taken or /// about which of its tags travel — the failure this consolidates away is the /// one where a second reading is added beside the first and only one of them /// is kept up to date. /// /// `{date}` is the capture date, not today's: a template naming exports by /// when the shutter fired is the reason the token exists. A file that recorded /// no capture time gets an empty string, never a substitute. /// The date and the carried record for one file's header — for the batch, /// and for a merge writing its composite's header from its first source. pub(crate) fn header_for_file(meta: &dr_decode::Metadata) -> (String, dr_export::SourceMetadata) { from_header(meta) } fn from_header(meta: &dr_decode::Metadata) -> (String, dr_export::SourceMetadata) { let date = meta .captured_at .map(crate::library_ui::format_date) .unwrap_or_default(); (date, carried_metadata(meta)) } /// TRACES: FR-EXP-8 /// What an export is allowed to carry from the file it was decoded from. /// /// Field by field rather than a conversion trait, and that is the point: /// `dr_export::SourceMetadata` is an allowlist, so a tag newly parsed by /// `dr-decode` reaches an exported file only when somebody adds a line here /// and thereby decides, in writing, that it may leave the machine. The /// location travels — `dr-export` is where the stripping decision is taken, /// once, from the settings, and duplicating it here would give two places to /// disagree. fn carried_metadata(meta: &dr_decode::Metadata) -> dr_export::SourceMetadata { dr_export::SourceMetadata { make: meta.make.clone(), model: meta.model.clone(), lens: meta.lens.clone(), shutter: meta.shutter, aperture: meta.aperture, iso: meta.iso, focal_length: meta.focal_length, captured_at: meta.captured_at, captured_offset: meta.captured_offset, artist: meta.artist.clone(), copyright: meta.copyright.clone(), location: meta.location, } } /// One rendered photograph on its way to a file: the name it will be written /// under, the name it came from, the pixels, and whatever metadata travelled /// with them. type RenderedItem = ( String, String, dr_export::Frame, Option, ); /// TRACES: FR-RAW-2 | FR-EXP-8 /// Read a fetched photograph's header and open it for rendering, through /// whichever decoder the batch was given. /// /// Split from [`render_from_library`] so the half that decodes can be run /// without a server: the stub-decoder test drives it directly. pub(crate) fn open_for_export( gpu: &dr_gpu::GpuContext, decoder: &dyn dr_decode::Decoder, bytes: &[u8], ) -> Result<(dr_decode::Metadata, crate::DevelopSession), String> { let meta = decoder.metadata(bytes).unwrap_or_default(); let session = crate::open_session(gpu, decoder, bytes, &meta)?; Ok((meta, session)) } /// Fetch a photograph, apply its stored edit, and render it at full size. /// /// `None` means the work was cancelled, which is not a failure and has no /// error to report — hence the `Option` outside the `Result`. fn render_from_library( request: &BatchRequest, path: &str, cache: Option, cancel: &Cancel, ) -> Option> { let Some(conn) = request.conn.clone() else { return Some(Err(ItemError::Fetch("no library is open".into()))); }; let Some(gpu) = request.gpu.as_ref() else { return Some(Err(ItemError::Open( "this build found no GPU adapter".into(), ))); }; // Both started before either is waited on, exactly as opening an image in // develop does: the sidecar is three orders of magnitude smaller than the // RAW, so it costs nothing to have in hand by the time there is a session // to apply it to. let sidecar_rx = crate::library::spawn_sidecar_fetch( conn.clone(), path.to_string(), request.sidecar_cache.clone(), request.offline, ); let bytes_rx = crate::library::spawn_full_fetch(conn, path.to_string(), cache); let bytes = match wait_for(&bytes_rx, cancel) { Waited::Got(Ok(bytes)) => bytes, Waited::Got(Err(e)) => return Some(Err(ItemError::Fetch(e.message))), Waited::Cancelled => return None, Waited::Silent => { return Some(Err(ItemError::Fetch( "the download ended without answering".into(), ))) } }; let sidecar = match wait_for(&sidecar_rx, cancel) { Waited::Got(sidecar) => sidecar, Waited::Cancelled => return None, // An unedited photograph has no sidecar and a fetch that died looks // identical from here. Exporting at defaults is what opening it would // do, and refusing the image over a missing edit it may never have had // would fail the common case. Waited::Silent => None, }; let (meta, mut session) = match open_for_export(gpu, request.decoder, &bytes) { Ok(opened) => opened, Err(e) => return Some(Err(ItemError::Open(e))), }; let (date, carried) = from_header(&meta); // TRACES: FR-CAT-8 // The stored edit. FR-EXP-9 is about resolution; this is the other half of // exporting what the user actually has, and a batch that skipped it would // write out three hundred unedited frames without saying so. if let Some(version) = sidecar .as_ref() .and_then(dr_pipeline::Sidecar::default_version) { session.apply_version(version); } // The last cheap place to stop. Everything past here is a full-resolution // render and an encode with no interior stopping point — see [`CANCEL_POLL`]. if cancel.is_cancelled() { return None; } let frame = match session.render_for_export(request.settings.colour_space) { Ok(frame) => frame, Err(e) => return Some(Err(ItemError::Render(e))), }; let stem = Path::new(path) .file_stem() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_else(|| "export".into()); // TRACES: FR-EXP-8 // `meta` was read at the top of this function — for the orientation the // session opens with as much as for this — and carrying it on to the // encoder is what puts the camera, the lens and the rights statement into // the exported file. What is *dropped* from it is decided in `dr-export` // from the settings, not here. Some(Ok((stem, date, frame, Some(carried)))) } /// Name, encode and write one rendered frame. /// /// The tail every source shares, whoever rendered it. fn place_frame( request: &BatchRequest, stem: &str, date: &str, sequence: u32, frame: &dr_export::Frame, // TRACES: FR-EXP-8 // What the source file said about itself, or `None` where the caller has // nothing to say. Handed straight through: every decision about what of it // reaches the file is taken inside `dr-export`, from the settings. source: Option<&dr_export::SourceMetadata>, issued: &mut HashSet, ) -> Result { // The size is resolved before the name because `{dimensions}` is one of the // tokens a template can carry. let (width, height) = dr_export::target_size( frame.width, frame.height, request.settings.sizing, request.settings.allow_upscaling, ); let ctx = NameContext { source_stem: stem, sequence, date, width, height, preset: "", }; let name = resolve_batch_name(&request.settings, &request.remote_names, &ctx, issued) .ok_or(ItemError::NameTaken)?; let encoded = dr_export::export(frame, &request.settings, name, source)?; place( &encoded, request.settings.target, &request.settings.destination, &request.outbox, request.settings.collision, ) .map_err(ItemError::Place) } /// The name this export takes, avoiding both what was in the destination and /// what this run has already written. /// /// Those are two different collisions and they deserve two different answers. /// [`dr_types::CollisionPolicy`] is the user's answer to "a file of this name /// was already there", and Overwrite is a perfectly reasonable one. It is not /// an answer to "the frame I exported four seconds ago was also called this": /// two photographs of the same stem in different folders are ordinary in a /// library, and a batch that quietly handed back fewer files than images — /// having destroyed its own output — is not something anybody asked for. So a /// name this run issued is always stepped past, whatever the policy says about /// the folder. fn resolve_batch_name( settings: &ExportSettings, remote_names: &HashSet, ctx: &NameContext<'_>, issued: &mut HashSet, ) -> Option { let dir = PathBuf::from(&settings.destination); let taken = |name: &str| -> bool { match settings.target { #[cfg(target_os = "android")] dr_types::ExportTarget::Device if settings.destination.starts_with("content://") => { crate::saf::exists(&settings.destination, name) } dr_types::ExportTarget::Device => dir.join(name).exists(), // A queued export cannot see the server, and may never be able to. // What the album records stands in for it here, and the upload // checks the server itself — see [`send`]. dr_types::ExportTarget::Remote => remote_names.contains(name), } }; let name = dr_export::resolve_name( &settings.filename_template, ctx, settings.format, settings.collision, &|name| taken(name) || issued.contains(name), )?; // Only reachable under Overwrite, which hands back the taken name by // design. Skip and Increment have already been through the closure above. let name = if issued.contains(&name) { step_past(&name, &|candidate| { taken(candidate) || issued.contains(candidate) })? } else { name }; issued.insert(name.clone()); Some(name) } /// `photo.jpg` → `photo-1.jpg`, and on until the name is free. fn step_past(name: &str, taken: &dyn Fn(&str) -> bool) -> Option { let path = Path::new(name); let stem = path .file_stem() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_else(|| "export".into()); let ext = path .extension() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_default(); // Bounded for the same reason `resolve_name`'s own search is: a destination // that reports every name as taken has to fail rather than spin. (1..10_000) .map(|n| { if ext.is_empty() { format!("{stem}-{n}") } else { format!("{stem}-{n}.{ext}") } }) .find(|candidate| !taken(candidate)) } /// What waiting on another worker produced. /// /// Three answers rather than an `Option`, because "the user cancelled" and "the /// worker died without a word" have the same shape and want opposite treatment: /// one ends the batch quietly, the other is a failure belonging to one image. enum Waited { Got(T), Cancelled, Silent, } /// Block on another worker's answer without going deaf to cancellation. fn wait_for(rx: &Receiver, cancel: &Cancel) -> Waited { loop { match rx.recv_timeout(CANCEL_POLL) { Ok(value) => return Waited::Got(value), Err(RecvTimeoutError::Timeout) => { if cancel.is_cancelled() { return Waited::Cancelled; } } Err(RecvTimeoutError::Disconnected) => return Waited::Silent, } } } /// Which status line a batch writes its commentary to. /// /// One drain serves both entry points rather than two near-identical timers, so /// it has to be told where the words go: develop has its own line beside the /// export button, and the grid has the header's. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Reporting { Develop, Library, } /// TRACES: FR-EXP-7 | NFR-P9 /// Drain a batch's progress on the UI thread. /// /// `slot` holds the timer, and dropping what was there before is what stops a /// superseded batch's drain — and, through [`crate::activity::Activity`]'s /// `Drop`, takes its row with it. pub fn drain_batch( weak: slint::Weak, activity: &Rc, slot: &Rc>>, rx: Receiver, total: usize, reporting: Reporting, // Run when the batch finishes having written something, with each file // it wrote and the image that file came from. A queued export is finished // on disk but not where the user asked for it, and waiting for the next // sync pass to notice reads — correctly — as an export that did not // upload; the files are what an album records (FR-EXP-10). on_exported: impl Fn(Vec<(dr_types::ImageId, String)>) + 'static, ) { let job = activity.begin( crate::activity::Kind::Export, if total == 1 { "Exporting".to_string() } else { format!("Exporting {total} images") }, ); job.total(total); // The first failure, kept for the summary. Only the first: a run where // every image failed for the same reason should say that reason once, and // the rest are in the log. let mut first_failure: Option = None; let mut written: Vec<(dr_types::ImageId, String)> = Vec::new(); let timer = slint::Timer::default(); let held = slot.clone(); timer.start(slint::TimerMode::Repeated, DRAIN_INTERVAL, move || { let Some(w) = weak.upgrade() else { return }; loop { let message = match rx.try_recv() { Ok(m) => m, Err(std::sync::mpsc::TryRecvError::Empty) => return, Err(std::sync::mpsc::TryRecvError::Disconnected) => { // A worker that died without a word must not leave the // button saying "Cancel export" for the rest of the session. job.fail("ended unexpectedly"); settle(&w, reporting, "Export ended unexpectedly"); stop_timer(&held); return; } }; match message { BatchMessage::Started { done, total, name } => { job.progress(done, total); job.detail(name.clone()); report( &w, reporting, &format!("Exporting {name} ({}/{total})", done + 1), ); } BatchMessage::Item { name, image, outcome, } => match outcome { Ok(placed) => { log::info!("{name}: {}", placed.describe()); if let Some(image) = image { written.push((image, placed.file_name())); } } Err(e) => { // Every failure is logged, not only the first: the // summary is a sentence and this is the record of which // photographs it is about (NFR-ARCH-4). log::warn!("exporting {name}: {e}"); first_failure.get_or_insert_with(|| format!("{name}: {e}")); } }, BatchMessage::Finished { exported, failed, cancelled, } => { let (text, is_failure) = summarise(exported, failed, cancelled, first_failure.as_deref()); if is_failure { job.fail(text.clone()); } else { job.finish(text.clone()); } settle(&w, reporting, &text); if exported > 0 { on_exported(std::mem::take(&mut written)); } stop_timer(&held); return; } } } }); *slot.borrow_mut() = Some(timer); } /// How a finished batch reads, and whether it counts as a failure. /// /// A cancelled run is never a failure however little it exported: the user /// stopped it, and a red row telling them so is the application arguing. A run /// with even one failure is, and stays in the list until it is cleared — that /// is the row somebody came looking for. fn summarise( exported: usize, failed: usize, cancelled: bool, first_failure: Option<&str>, ) -> (String, bool) { if cancelled { return (format!("Cancelled after {exported}"), false); } if failed == 0 { return (format!("Exported {exported}"), false); } let detail = first_failure.unwrap_or("see the log"); ( format!("Exported {exported}, {failed} failed — {detail}"), true, ) } /// Say what the batch is doing on the line the caller came from. fn report(window: &AppWindow, reporting: Reporting, text: &str) { match reporting { Reporting::Develop => window.set_export_status(text.into()), Reporting::Library => window.global::().set_library_status(text.into()), } } /// Report, and put the buttons back. fn settle(window: &AppWindow, reporting: Reporting, text: &str) { report(window, reporting, text); match reporting { Reporting::Develop => window.set_export_busy(false), Reporting::Library => window.global::().set_library_exporting(false), } } fn stop_timer(slot: &Rc>>) { if let Some(timer) = slot.borrow().as_ref() { timer.stop(); } } #[cfg(test)] 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, &mut Listings::new(), ) .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, &mut Listings::new()) .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, &mut Listings::new(), ) .await .unwrap(); assert_eq!(server.lists.load(Ordering::SeqCst), 0); let _ = std::fs::remove_dir_all(&s.dir); } /// An album export queued under `policy`, the album recording it, and a /// server that already holds a file of that name. async fn queued_over_a_taken_name( tag: &str, policy: CollisionPolicy, ) -> (Staged, Assigning, dr_catalog::albums::AlbumId) { let s = staged(tag, "PhotosRaw"); let catalog = dr_catalog::Catalog::open(&s.library.catalog).unwrap(); let album = dr_catalog::albums::create( catalog.connection(), "Web", &dr_catalog::albums::Place::Server("Shared/Web".into()), ) .unwrap(); dr_catalog::albums::record_exports( catalog.connection(), album, &[(dr_types::ImageId(s.image as u64), "a.jpg".into())], ) .unwrap(); let server = Assigning::default(); use dr_sync::RemoteBackend as _; server .put( &RemotePath::new("Shared/Web/a.jpg"), b"theirs".to_vec(), None, ) .await .unwrap(); place( &encoded("a.jpg", b"ours"), ExportTarget::Remote, "/Shared/Web", &s.dir.join("outbox"), policy, ) .unwrap(); (s, server, album) } fn queued_named(s: &Staged, name: &str) -> Vec { pending(&s.dir.join("outbox")) .into_iter() .filter(|p| p.name == name) .collect() } fn held(server: &Assigning, path: &str) -> Option> { server .files .lock() .unwrap() .get(path) .map(|(_, body)| body.clone()) } #[tokio::test] async fn an_increment_export_steps_past_a_name_the_server_holds() { let (s, server, album) = queued_over_a_taken_name("increment", CollisionPolicy::Increment).await; // A second export of the same name, queued before either uploaded. place( &encoded("a.jpg", b"ours too"), ExportTarget::Remote, "/Shared/Web", &s.dir.join("outbox"), CollisionPolicy::Increment, ) .unwrap(); let mut listings = Listings::new(); for entry in queued_named(&s, "a.jpg") { let sent = send(&server, &s.library, "PhotosRaw", &entry, &mut listings) .await .unwrap(); assert_eq!(sent, Sent::Uploaded); } assert_eq!(held(&server, "Shared/Web/a.jpg").unwrap(), b"theirs"); let mut ours = vec![ held(&server, "Shared/Web/a-1.jpg").unwrap(), held(&server, "Shared/Web/a-2.jpg").unwrap(), ]; ours.sort(); assert_eq!(ours, vec![b"ours".to_vec(), b"ours too".to_vec()]); assert_eq!( server.lists.load(Ordering::SeqCst), 1, "one listing for the folder, however many files go into it" ); // The album's row followed its file to the name it was given. let catalog = dr_catalog::Catalog::open(&s.library.catalog).unwrap(); let names = dr_catalog::albums::file_names(catalog.connection(), album).unwrap(); assert!(!names.contains("a.jpg"), "{names:?}"); assert!( queued_named(&s, "a.jpg").is_empty(), "the outbox is cleared" ); let _ = std::fs::remove_dir_all(&s.dir); } #[tokio::test] async fn a_skip_export_leaves_a_name_the_server_holds() { let (s, server, _) = queued_over_a_taken_name("skip", CollisionPolicy::Skip).await; let entry = queued_named(&s, "a.jpg").pop().unwrap(); let sent = send( &server, &s.library, "PhotosRaw", &entry, &mut Listings::new(), ) .await .unwrap(); assert_eq!(sent, Sent::Skipped); assert_eq!(held(&server, "Shared/Web/a.jpg").unwrap(), b"theirs"); assert_eq!(server.files.lock().unwrap().len(), 1); assert!( queued_named(&s, "a.jpg").is_empty(), "the outbox is cleared" ); let _ = std::fs::remove_dir_all(&s.dir); } #[tokio::test] async fn an_overwrite_export_replaces_a_name_the_server_holds() { let (s, server, _) = queued_over_a_taken_name("overwrite", CollisionPolicy::Overwrite).await; let entry = queued_named(&s, "a.jpg").pop().unwrap(); send( &server, &s.library, "PhotosRaw", &entry, &mut Listings::new(), ) .await .unwrap(); assert_eq!(held(&server, "Shared/Web/a.jpg").unwrap(), b"ours"); 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(), bytes: bytes.to_vec(), width: 4, height: 4, } } fn tmp() -> PathBuf { let dir = std::env::temp_dir().join(format!( "dr-outbox-test-{}-{:?}", std::process::id(), std::thread::current().id() )); let _ = std::fs::remove_dir_all(&dir); std::fs::create_dir_all(&dir).unwrap(); dir } #[test] fn a_device_export_writes_the_file() { let dir = tmp(); let target = dir.join("exports"); let placed = place( &encoded("a.jpg", b"hello"), ExportTarget::Device, target.to_str().unwrap(), &dir.join("outbox"), CollisionPolicy::Increment, ) .unwrap(); assert_eq!(placed, Placed::Device(target.join("a.jpg"))); assert_eq!(std::fs::read(target.join("a.jpg")).unwrap(), b"hello"); } #[test] fn a_device_export_creates_a_folder_that_is_not_there() { // Exporting into a folder the user typed but has not made is the // common case, not an error. let dir = tmp(); let target = dir.join("deep/nested/exports"); assert!(place( &encoded("a.jpg", b"x"), ExportTarget::Device, target.to_str().unwrap(), &dir, CollisionPolicy::Increment, ) .is_ok()); assert!(target.join("a.jpg").exists()); } #[test] fn a_device_export_with_no_folder_says_so() { // Rather than writing to the process's working directory, which is // wherever the app happened to be launched from. let dir = tmp(); let err = place( &encoded("a.jpg", b"x"), ExportTarget::Device, " ", &dir, CollisionPolicy::Increment, ) .unwrap_err(); // Says what to do about it: the destination is an album now. assert!(err.contains("album"), "unhelpful message: {err}"); } #[test] fn a_remote_export_is_staged_rather_than_sent() { // The property the offline story rests on: the export is complete on // disk before any network call is attempted. let dir = tmp(); let outbox = dir.join("outbox"); let placed = place( &encoded("a.jpg", b"hello"), ExportTarget::Remote, "Exports/2026", &outbox, CollisionPolicy::Increment, ) .unwrap(); match placed { Placed::Queued { local, remote_dir, .. } => { assert_eq!(std::fs::read(&local).unwrap(), b"hello"); assert_eq!(remote_dir, "Exports/2026"); } other => panic!("expected a queued export, got {other:?}"), } } #[test] fn a_staged_export_is_found_again_with_its_destination() { // What survives a process death: the drain has to be able to // reconstruct where a file was going from the disk alone. let dir = tmp(); let outbox = dir.join("outbox"); place( &encoded("a.jpg", b"one"), ExportTarget::Remote, "Exports", &outbox, CollisionPolicy::Increment, ) .unwrap(); let queue = pending(&outbox); assert_eq!(queue.len(), 1); assert_eq!(queue[0].remote_dir, "Exports"); assert_eq!(queue[0].name, "a.jpg"); } #[test] fn two_exports_of_the_same_name_both_survive_the_outbox() { // Both were asked for and neither has uploaded, so the second must // not overwrite the first while it waits. let dir = tmp(); let outbox = dir.join("outbox"); place( &encoded("a.jpg", b"one"), ExportTarget::Remote, "E", &outbox, CollisionPolicy::Increment, ) .unwrap(); place( &encoded("a.jpg", b"two"), ExportTarget::Remote, "E", &outbox, CollisionPolicy::Increment, ) .unwrap(); let queue = pending(&outbox); assert_eq!(queue.len(), 2); // Both still claim the name they should arrive under; only the local // staging name differs. assert!(queue.iter().all(|p| p.name == "a.jpg")); assert_ne!(queue[0].local, queue[1].local); } #[test] fn a_remote_batch_names_around_the_album_and_the_outbox() { let dir = tmp(); let outbox = dir.join("outbox"); // Queued by an earlier batch, not uploaded yet. place( &encoded("IMG_0001-1.png", b"x"), ExportTarget::Remote, "/Shared/Web", &outbox, CollisionPolicy::Increment, ) .unwrap(); // Elsewhere, so in nobody's way. place( &encoded("IMG_0001-2.png", b"x"), ExportTarget::Remote, "Shared/Web", &outbox, CollisionPolicy::Increment, ) .unwrap(); let settings = ExportSettings { format: dr_types::ExportFormat::Png, target: ExportTarget::Remote, destination: "/Shared/Web".into(), collision: CollisionPolicy::Increment, ..Default::default() }; // What the album records of an earlier export. let mut known: HashSet = ["IMG_0001.png".to_string()].into(); known.extend(queued_for(&outbox, &settings.destination)); let name = resolve_batch_name( &settings, &known, &NameContext { source_stem: "IMG_0001", sequence: 1, ..Default::default() }, &mut HashSet::new(), ); assert_eq!(name.as_deref(), Some("IMG_0001-2.png")); } #[test] fn a_record_carries_its_policy_and_an_old_one_has_none() { let dir = tmp(); let outbox = dir.join("outbox"); place( &encoded("a.jpg", b"x"), ExportTarget::Remote, "E", &outbox, CollisionPolicy::Skip, ) .unwrap(); let queue = pending(&outbox); assert_eq!(queue[0].collision, Some(CollisionPolicy::Skip)); assert!(!queue[0].account); std::fs::write(outbox.join("old.jpg"), b"x").unwrap(); std::fs::write(outbox.join("old.jpg.dest"), "E\nold.jpg\naccount\n").unwrap(); let old = pending(&outbox) .into_iter() .find(|p| p.name == "old.jpg") .unwrap(); assert_eq!(old.collision, None); assert!(old.account); } #[test] fn a_payload_with_no_record_is_ignored() { // The window a kill between the two writes leaves behind. It must not // become an upload to nowhere. let dir = tmp(); let outbox = dir.join("outbox"); std::fs::create_dir_all(&outbox).unwrap(); std::fs::write(outbox.join("orphan.jpg"), b"x").unwrap(); assert!(pending(&outbox).is_empty()); } #[test] fn a_record_with_no_payload_is_ignored() { let dir = tmp(); let outbox = dir.join("outbox"); std::fs::create_dir_all(&outbox).unwrap(); std::fs::write(outbox.join("ghost.jpg.dest"), "E\nghost.jpg\n").unwrap(); assert!(pending(&outbox).is_empty()); } #[test] fn an_empty_outbox_is_not_an_error() { // Called on every sync pass, including before anything is exported // and on a device where the directory has never been created. assert!(pending(Path::new("/nonexistent/darkroom/outbox")).is_empty()); assert_eq!(pending_count(Path::new("/nonexistent/darkroom/outbox")), 0); } #[test] fn the_remote_path_joins_root_folder_and_name() { let entry = Pending { local: PathBuf::from("/tmp/a.jpg"), remote_dir: "Exports/2026".into(), name: "a.jpg".into(), account: false, collision: None, }; assert_eq!( entry.remote_path("Photos").as_str(), "Photos/Exports/2026/a.jpg" ); assert_eq!( entry.remote_folder("Photos").as_str(), "Photos/Exports/2026" ); } #[test] fn an_empty_folder_exports_to_the_library_root() { // "Ask each time" is not set here — an empty remote folder means the // root, and it must not produce a double slash the server rejects. let entry = Pending { local: PathBuf::from("/tmp/a.jpg"), remote_dir: String::new(), name: "a.jpg".into(), account: false, collision: None, }; assert_eq!(entry.remote_path("Photos").as_str(), "Photos/a.jpg"); } #[test] fn stray_slashes_do_not_produce_an_unusable_path() { // The folder comes from a picker or a text field, and either can hand // over a leading or trailing slash. let entry = Pending { local: PathBuf::from("/tmp/a.jpg"), remote_dir: "/Exports/".into(), name: "a.jpg".into(), account: false, collision: None, }; assert_eq!( entry.remote_path("/Photos/").as_str(), "Photos/Exports/a.jpg" ); } #[test] fn an_album_on_the_server_is_reached_from_the_account_not_the_library() { // FR-EXP-10: an album lives outside the library, so its folder is // spelled from the account root and the library root is not put in // front of it — through the record on disk, as a drain after a // restart would read it. let dir = std::env::temp_dir().join(format!( "dr-export-album-outbox-{}-{:?}", std::process::id(), std::thread::current().id() )); let _ = std::fs::remove_dir_all(&dir); place( &encoded("a.jpg", b"x"), ExportTarget::Remote, "/Shared/Web", &dir, CollisionPolicy::Increment, ) .unwrap(); let entry = pending(&dir).pop().expect("one staged export"); assert!(entry.account); assert_eq!(entry.remote_path("Photos").as_str(), "Shared/Web/a.jpg"); assert_eq!(entry.remote_folder("Photos").as_str(), "Shared/Web"); let _ = std::fs::remove_dir_all(&dir); } #[test] fn the_status_line_never_claims_an_upload_that_has_not_happened() { // A queued export is real and finished, but it is not on the server, // and saying so before it is would be a promise the app cannot keep. let queued = Placed::Queued { local: PathBuf::from("/tmp/a.jpg"), remote_dir: "Exports".into(), name: "a.jpg".into(), }; let text = queued.describe(); assert!(text.contains("Queued"), "{text}"); assert!(!text.contains("Exported"), "{text}"); assert!(Placed::Device(PathBuf::from("/tmp/a.jpg")) .describe() .contains("Exported")); } // ----------------------------------------------------------------------- // The batch (FR-EXP-7) // ----------------------------------------------------------------------- /// Settings that write PNGs into `dir`, so a test can look at what landed. fn to_folder(dir: &Path) -> ExportSettings { ExportSettings { format: dr_types::ExportFormat::Png, target: ExportTarget::Device, destination: dir.to_string_lossy().into_owned(), ..Default::default() } } fn request(settings: ExportSettings, sources: Vec) -> BatchRequest { BatchRequest { sources, images: Vec::new(), conn: None, settings, remote_names: HashSet::new(), outbox: std::env::temp_dir().join("dr-batch-test-outbox"), sidecar_cache: std::env::temp_dir().join("dr-batch-test-sidecars"), offline: true, gpu: None, decoder: dr_decode::default(), } } /// A frame of flat pixels — enough for the encoder, cheap for a test. fn frame(width: u32, height: u32) -> dr_export::Frame { dr_export::Frame::new(width, height, vec![128; (width * height * 4) as usize]) .expect("well-formed") } fn drive(request: BatchRequest, cancel: Cancel) -> Vec { let (tx, rx) = std::sync::mpsc::channel(); run(request, &cancel, &tx); drop(tx); rx.into_iter().collect() } fn finished(messages: &[BatchMessage]) -> (usize, usize, bool) { messages .iter() .find_map(|m| match m { BatchMessage::Finished { exported, failed, cancelled, } => Some((*exported, *failed, *cancelled)), _ => None, }) .expect("a batch always says it has finished") } // ----------------------------------------------------------------------- // What an export says about the photograph (FR-EXP-8) // ----------------------------------------------------------------------- /// A header of the kind a camera writes, GPS included. /// /// The capture time is 2013-06-28 23:32:54 UTC, which `dr-export`'s own /// round-trip test uses for the same reason: it is a date nothing else in /// the pipeline could have produced by accident. fn header() -> dr_decode::Metadata { dr_decode::Metadata { make: Some("Canon".into()), model: Some("EOS 6D".into()), lens: Some("EF 35mm f/2 IS USM".into()), shutter: Some(1.0 / 250.0), aperture: Some(2.8), iso: Some(400), focal_length: Some(35.0), captured_at: Some(1_372_462_374), captured_offset: Some(120), artist: Some("A Photographer".into()), copyright: Some("(c) A Photographer".into()), location: dr_types::Location::new(48.8582, 2.2945, Some(35.0)), ..Default::default() } } /// Settings that write JPEGs into `dir`. /// /// JPEG rather than the PNG `to_folder` above uses. `dr-export` writes the /// same EXIF block into a PNG's `eXIf` chunk, but `dr-decode` reads a /// header out of a JPEG, so this is the format where the assertion can be /// made against what a reader actually finds rather than against the bytes /// this crate just wrote. fn to_folder_as_jpeg(dir: &Path) -> ExportSettings { ExportSettings { format: dr_types::ExportFormat::Jpeg, target: ExportTarget::Device, destination: dir.to_string_lossy().into_owned(), ..Default::default() } } /// Run one already-rendered frame through the batch, as the develop button /// does, and hand back what landed. fn develop_export( settings: &ExportSettings, header: Option, stem: &str, ) -> Vec { let messages = drive( request( settings.clone(), vec![Source::Rendered { stem: stem.into(), header, frame: frame(16, 12), }], ), Cancel::default(), ); assert_eq!(finished(&messages), (1, 0, false), "the export itself"); let written = messages .iter() .find_map(|m| match m { BatchMessage::Item { outcome: Ok(Placed::Device(path)), .. } => Some(path.clone()), _ => None, }) .expect("a device export names the file it wrote"); std::fs::read(written).expect("the file the batch says it wrote") } #[test] fn a_develop_export_writes_the_file_a_library_export_would_have() { // The two buttons on the same photograph. The library path renders the // frame itself and then hands the encoder `(date, carried_metadata)`; // the develop path hands over a frame already rendered. Everything // after that is shared, so if the develop path's file is byte for byte // what the library path's arguments produce, the two cannot differ in // what they disclose — which is the whole of this bug. let dir = tmp(); let out = dir.join("exports"); let settings = to_folder_as_jpeg(&out); let written = develop_export(&settings, Some(header()), "IMG_0001"); let (_, carried) = from_header(&header()); let library = dr_export::export( &frame(16, 12), &settings, "IMG_0001.jpg".into(), Some(&carried), ) .expect("the same frame and the same settings"); assert_eq!( written, library.bytes, "a develop export and a library export of one photograph differ" ); // And it is a header a reader can find, not merely bytes that match: // an empty block on both sides would satisfy the comparison above. let read_back = dr_decode::metadata(&written).expect("our own JPEG"); assert_eq!(read_back.make.as_deref(), Some("Canon")); assert_eq!(read_back.model.as_deref(), Some("EOS 6D")); assert_eq!(read_back.captured_at, header().captured_at); assert_eq!(read_back.copyright, header().copyright); } #[test] fn the_date_token_on_a_develop_export_is_when_the_shutter_fired() { // Not today's date. A photographer filing an export of a negative // scanned this morning wants 1978 in the name, and the library path // has always given them that. let dir = tmp(); let out = dir.join("exports"); let mut settings = to_folder_as_jpeg(&out); settings.filename_template = "{date}_{name}".into(); develop_export(&settings, Some(header()), "IMG_0001"); assert!( out.join("2013-06-28_IMG_0001.jpg").exists(), "{:?}", std::fs::read_dir(&out).unwrap().flatten().count() ); } #[test] fn a_develop_export_obeys_the_location_setting_the_same_way() { // The stripping decision is `dr-export`'s and is taken from the // settings; this path must reach it rather than pre-filtering, or the // photographer who deliberately turned stripping off would find their // coordinates dropped anyway on one of the two buttons. let dir = tmp(); let out = dir.join("exports"); let mut stripped = to_folder_as_jpeg(&out); stripped.strip_location = true; let bytes = develop_export(&stripped, Some(header()), "IMG_0001"); assert_eq!( dr_decode::metadata(&bytes).expect("our own JPEG").location, None, "coordinates survived an export that was told to strip them" ); let mut kept = to_folder_as_jpeg(&out); kept.strip_location = false; let bytes = develop_export(&kept, Some(header()), "IMG_0002"); let read_back = dr_decode::metadata(&bytes) .expect("our own JPEG") .location .expect("the position the settings asked to keep"); // Through degrees, minutes and seconds and back, so exactly is the // wrong word; a metre is about 1e-5 degrees. assert!((read_back.latitude - 48.8582).abs() < 1e-5, "{read_back:?}"); assert!((read_back.longitude - 2.2945).abs() < 1e-5, "{read_back:?}"); } #[test] fn a_photograph_with_no_header_exports_and_invents_no_date() { // `{date}` has nothing to say and says nothing. The alternative — // quietly substituting today — would date every scan and every // headerless JPEG to the afternoon it was exported, and the filename // is exactly where that lie would be hardest to notice. let dir = tmp(); let out = dir.join("exports"); let mut settings = to_folder_as_jpeg(&out); settings.filename_template = "{date}{name}".into(); // A file with no EXIF block, which is what the develop path really // hands over for one: `load_bytes` reads the header as // `Metadata::default()` and the session remembers that, exactly as the // library path does for the same file. let bytes = develop_export(&settings, Some(dr_decode::Metadata::default()), "IMG_0001"); assert!(out.join("IMG_0001.jpg").exists()); assert_eq!( dr_decode::metadata(&bytes) .expect("our own JPEG") .captured_at, None ); // And a session with no header at all — the fixture case, and a // decoder that returned nothing. `dr-export` writes no EXIF segment // rather than an empty one, which `dr-decode` reports as an error // rather than as a blank header; both readings are "no date", and the // point of the test is that neither is today's. let bytes = develop_export(&settings, None, "IMG_0002"); assert!(out.join("IMG_0002.jpg").exists()); assert_eq!( dr_decode::metadata(&bytes).ok().and_then(|m| m.captured_at), None ); } #[test] fn a_rendered_frame_reaches_the_destination_folder() { // The whole tail of the batch — name, size, sharpen, encode, write — // with no GPU and no network, which is what makes it testable at all. let dir = tmp(); let out = dir.join("exports"); let messages = drive( request( to_folder(&out), vec![Source::Rendered { stem: "IMG_0001".into(), header: None, frame: frame(16, 12), }], ), Cancel::default(), ); assert_eq!(finished(&messages), (1, 0, false)); assert!(out.join("IMG_0001.png").exists()); } #[test] fn two_images_of_one_name_both_survive_the_batch() { // Two folders in a library holding an IMG_0001 each is ordinary, and a // batch that wrote one file for two photographs would destroy work // without saying anything. Overwrite is set deliberately: it is the // user's answer about the *folder*, not about this run's own output. let dir = tmp(); let out = dir.join("exports"); let mut settings = to_folder(&out); settings.collision = dr_types::CollisionPolicy::Overwrite; let messages = drive( request( settings, vec![ Source::Rendered { stem: "IMG_0001".into(), header: None, frame: frame(16, 12), }, Source::Rendered { stem: "IMG_0001".into(), header: None, frame: frame(16, 12), }, ], ), Cancel::default(), ); assert_eq!(finished(&messages), (2, 0, false)); assert!(out.join("IMG_0001.png").exists()); assert!(out.join("IMG_0001-1.png").exists()); } #[test] fn a_name_that_was_already_there_still_obeys_the_collision_setting() { // The other half of the rule above: a file that existed *before* the // batch is exactly what the policy is about, and Overwrite must still // mean overwrite or the setting would do nothing. let dir = tmp(); let out = dir.join("exports"); std::fs::create_dir_all(&out).expect("temp dir"); std::fs::write(out.join("IMG_0001.png"), b"older").expect("seed"); let mut settings = to_folder(&out); settings.collision = dr_types::CollisionPolicy::Overwrite; let mut issued = HashSet::new(); let name = resolve_batch_name( &settings, &HashSet::new(), &NameContext { source_stem: "IMG_0001", sequence: 1, ..Default::default() }, &mut issued, ); assert_eq!(name.as_deref(), Some("IMG_0001.png")); } #[test] fn a_skipped_name_is_reported_rather_than_silently_dropped() { // Skip is a legitimate answer, but the image still produced no file — // and a batch claiming to have exported it would be lying. let dir = tmp(); let out = dir.join("exports"); std::fs::create_dir_all(&out).expect("temp dir"); std::fs::write(out.join("IMG_0001.png"), b"older").expect("seed"); let mut settings = to_folder(&out); settings.collision = dr_types::CollisionPolicy::Skip; let messages = drive( request( settings, vec![Source::Rendered { stem: "IMG_0001".into(), header: None, frame: frame(8, 8), }], ), Cancel::default(), ); assert_eq!(finished(&messages), (0, 1, false)); assert!(matches!( messages.iter().find_map(|m| match m { BatchMessage::Item { outcome, .. } => Some(outcome), _ => None, }), Some(Err(ItemError::NameTaken)) )); assert_eq!(std::fs::read(out.join("IMG_0001.png")).unwrap(), b"older"); } #[test] fn one_failure_does_not_abandon_the_rest_of_the_batch() { // FR-EXP-7's central promise. The first source cannot be fetched — no // library is open — and the two after it must still be attempted and // still be counted. let dir = tmp(); let out = dir.join("exports"); let messages = drive( request( to_folder(&out), vec![ Source::Library { path: "Photos/broken.CR2".into(), cache: None, }, Source::Rendered { stem: "good-a".into(), header: None, frame: frame(8, 8), }, Source::Rendered { stem: "good-b".into(), header: None, frame: frame(8, 8), }, ], ), Cancel::default(), ); assert_eq!(finished(&messages), (2, 1, false)); assert!(out.join("good-a.png").exists()); assert!(out.join("good-b.png").exists()); } #[test] fn every_image_gets_its_own_outcome() { // The report is per image, not one verdict for the run: two failures // for two different reasons have to arrive as two messages, or a // three-hundred-frame batch is unreportable (NFR-ARCH-4). let dir = tmp(); let messages = drive( request( to_folder(&dir.join("exports")), vec![ Source::Library { path: "Photos/a.CR2".into(), cache: None, }, Source::Rendered { stem: "b".into(), header: None, frame: frame(8, 8), }, ], ), Cancel::default(), ); let outcomes: Vec<&Result> = messages .iter() .filter_map(|m| match m { BatchMessage::Item { outcome, .. } => Some(outcome), _ => None, }) .collect(); assert_eq!(outcomes.len(), 2); assert!(outcomes[0].is_err()); assert!(outcomes[1].is_ok()); } #[test] fn a_cancelled_batch_stops_and_writes_nothing_more() { // TRACES: NFR-ARCH-3 // Cancelled before it began, which is the strongest form of the // property: not one file, and the run still reports itself finished // rather than leaving the interface waiting for a message. let dir = tmp(); let out = dir.join("exports"); let cancel = Cancel::default(); cancel.cancel(); let messages = drive( request( to_folder(&out), vec![Source::Rendered { stem: "IMG_0001".into(), header: None, frame: frame(8, 8), }], ), cancel, ); assert_eq!(finished(&messages), (0, 0, true)); assert!(!out.join("IMG_0001.png").exists()); } #[test] fn a_cancelled_wait_gives_up_instead_of_blocking_for_ever() { // TRACES: NFR-ARCH-3 // The bound on cancelling a batch that is waiting on a download. With // a plain `recv` this test would hang, which is precisely the bug. let (tx, rx) = std::sync::mpsc::channel::(); let cancel = Cancel::default(); cancel.cancel(); assert!(matches!(wait_for(&rx, &cancel), Waited::Cancelled)); drop(tx); } #[test] fn a_worker_that_dies_is_not_mistaken_for_a_cancellation() { // The two look identical from a channel and mean opposite things: one // ends the batch, the other fails one image and moves on. let (tx, rx) = std::sync::mpsc::channel::(); drop(tx); assert!(matches!(wait_for(&rx, &Cancel::default()), Waited::Silent)); } #[test] fn a_cancelled_run_is_not_reported_as_a_failure() { // The user stopped it. A red row saying so is the application arguing // with something it was told to do. let (text, failed) = summarise(5, 0, true, None); assert!(!failed, "{text}"); assert!(text.contains('5'), "{text}"); } #[test] fn a_run_with_a_failure_keeps_its_row_and_names_one() { // This is the row somebody comes to the activity list to find, so it // has to survive being trimmed — and it has to say which photograph. let (text, failed) = summarise(38, 2, false, Some("IMG_0007.CR2: could not be rendered")); assert!(failed); assert!(text.contains("IMG_0007.CR2"), "{text}"); assert!(text.contains("38"), "{text}"); } }