//! Fetching thumbnails, previews and originals from the remote: the pin //! cache, dehydration, and the prefetcher that keeps the grid ahead of //! scrolling. use crate::executors::{self, Executor}; use crate::sidecar_cache::SidecarCache; use dr_catalog::Catalog; #[cfg(test)] use dr_sync::Account; use dr_sync::{Connection, RemoteId, RemotePath}; use std::path::PathBuf; use std::sync::mpsc::{Receiver, Sender}; use super::scan::now_secs; use super::sidecar::sidecar_path; /// Largest preview worth fetching whole. /// /// A located preview above this is skipped rather than transferred: past a few /// MB the saving over the full file stops justifying the wait, and a 256px /// thumbnail needs nothing like that much detail. pub(super) const MAX_PREVIEW_BYTES: u64 = 8 * 1024 * 1024; /// What the grid wants a thumbnail for. /// /// Carries the `oc:fileid` as well as the path, because that is what the /// shared store keys on — stable across a server-side move, and the same id /// every other client sees (FR-NC-5). #[derive(Debug, Clone)] pub struct ThumbnailRequest { pub row: usize, pub path: String, /// `None` where the scan found no stable id; such an image is fetched but /// not stored, since there is no durable key to store it under. pub file_id: Option, /// File length, needed to reject a preview range that points past the end /// of the file (NFR-SEC-1). pub size: u64, /// Catalog row, so EXIF read from the header can be written back. pub image_id: i64, /// Which resolution this cell needs, from how large it is drawn. A zoomed /// grid asks for the large class; a wall of small cells does not. /// /// Named apart from `size`, which is the file's length in bytes — the two /// are unrelated and confusing them would fetch the wrong thing. pub thumb_size: dr_thumbs::ThumbSize, /// Whether this image still needs its EXIF read. Where false the header is /// still fetched — the preview needs it — but nothing is parsed or written. pub needs_metadata: bool, /// Keep the preview at the resolution it was decoded at, ignoring /// `thumb_size`. /// /// For face indexing, which wants the pixels a thumbnail throws away: a /// face 2% across the frame is 5 px on a grid thumbnail and 120 px on the /// embedded preview, and 112 is what the embedder samples. Capped by /// [`FACE_SOURCE_EDGE`] rather than truly unbounded, because a 24 MP buffer /// converted to `f32` RGB is ~288 MB and several lanes hold one at once. pub full_resolution: bool, } /// Why a full fetch failed, keeping the one bit the UI cannot re-derive. /// /// The same reasoning as [`ScanFailure`]: the typed error cannot cross the /// channel, and "offline" versus "refused" decides whether develop shows /// "you are offline — this image is not stored locally" or a real error. #[derive(Debug)] pub struct FetchFailure { pub message: String, pub offline: bool, } impl FetchFailure { fn local(message: impl std::fmt::Display) -> Self { Self { message: message.to_string(), offline: false, } } } impl From for FetchFailure { fn from(e: dr_sync::RemoteError) -> Self { Self { offline: e.indicates_offline(), message: e.to_string(), } } } impl std::fmt::Display for FetchFailure { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.write_str(&self.message) } } /// TRACES: FR-NC-6a /// Progress from the pin worker. #[derive(Debug)] pub enum PinMessage { /// How many originals the pin still needs. Sent once, before any transfer. Planned { total: usize, }, /// One original landed. Stored { done: usize, }, /// The pin is fully downloaded. Done { stored: usize, bytes: u64, }, Failed { message: String, offline: bool, }, } /// TRACES: FR-NC-6a /// Download every original a pin has asked for. /// /// Whole files, deliberately: a pin exists so the photographs can be *edited* /// away from the server, and develop needs every photosite. This is the one /// place in the app that fetches originals in bulk, which is why FR-NC-6 /// makes it opt-in rather than something sync does on its own. /// /// Sequential rather than parallel. The lanes that make the thumbnail sweep /// fast are wrong here: these are tens of megabytes each, so concurrency buys /// little against a single connection's bandwidth and costs a great deal of /// memory — and it is the same contention that produced 423 Locked in the /// sweep. pub fn spawn_pin_fetch( conn: Connection, catalog_path: PathBuf, cache_dir: PathBuf, budget: dr_catalog::Budget, ) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "pin", move || { let (store, catalog) = match ( dr_catalog::Cache::open(&cache_dir, budget), Catalog::open(&catalog_path), ) { (Ok(s), Ok(c)) => (s, c), (Err(e), _) | (_, Err(e)) => { let _ = tx.send(PinMessage::Failed { message: e.to_string(), offline: false, }); return; } }; let pending = match store.pending_pins(catalog.connection()) { Ok(p) => p, Err(e) => { let _ = tx.send(PinMessage::Failed { message: e.to_string(), offline: false, }); return; } }; if tx .send(PinMessage::Planned { total: pending.len(), }) .is_err() { return; } if pending.is_empty() { let _ = tx.send(PinMessage::Done { stored: 0, bytes: 0, }); return; } let rt = match crate::net_runtime::build() { Ok(rt) => rt, Err(e) => { let _ = tx.send(PinMessage::Failed { message: e.to_string(), offline: false, }); return; } }; rt.block_on(async { let backend = match crate::remote::connect(&conn) { Ok(b) => b, Err(e) => { let _ = tx.send(PinMessage::Failed { message: e.to_string(), offline: false, }); return; } }; let mut stored = 0usize; let mut bytes_total = 0u64; for image in pending { let Some(source_ref) = source_ref_of(&catalog, image) else { // Catalogued and then removed while the pin was pending. continue; }; let id = RemoteId::Path(RemotePath::new(&source_ref)); // TRACES: FR-NC-6c // On a placeholder library "pin" means *keep it downloaded*, // not "make a second copy". The original materialises in the // library folder itself, so copying it under `originals/` // would hold every pinned photograph twice — and the copy // would be the half the budget could evict while the real disk // cost stayed. Only the bookkeeping is recorded, with no path, // so nothing here can ever delete a file inside a synced tree // (see `Cache::record_in_place`). if backend.capabilities().materialisation.can_materialise() { match backend.materialise(&id).await { Ok(_) => { let bytes = size_of(&catalog, image).unwrap_or(0); if let Err(e) = store.record_in_place( catalog.connection(), image, bytes, true, now_secs(), ) { log::warn!("recording pinned {source_ref}: {e}"); continue; } stored += 1; bytes_total += bytes; if tx.send(PinMessage::Stored { done: stored }).is_err() { return; } } Err(e) if e.indicates_offline() => { let _ = tx.send(PinMessage::Failed { message: e.to_string(), offline: true, }); return; } Err(e) => log::warn!("pinning {source_ref}: {e}"), } continue; } match backend.get(&id, None).await { Ok(bytes) => { // `pinned: true` — this is the population the budget // must never evict, which is the entire promise the // user made when they pinned the collection. if let Err(e) = store.store( catalog.connection(), image, &source_ref, &bytes, true, now_secs(), ) { log::warn!("storing pinned {source_ref}: {e}"); continue; } stored += 1; bytes_total += bytes.len() as u64; if tx.send(PinMessage::Stored { done: stored }).is_err() { return; } } Err(e) if e.indicates_offline() => { // Stop rather than failing each remaining file against // a dead connection. What was downloaded stays // downloaded, and `pending_pins` resumes from there. let _ = tx.send(PinMessage::Failed { message: e.to_string(), offline: true, }); return; } Err(e) => { // One unreadable file must not abandon the whole pin. log::warn!("pinning {source_ref}: {e}"); } } } let _ = tx.send(PinMessage::Done { stored, bytes: bytes_total, }); }); }); rx } /// TRACES: FR-NC-6c /// Hand a set of photographs back to the sync client, freeing their disk. /// /// The other half of pinning on a placeholder library. `Cache::release` drops /// the bookkeeping and — correctly — deletes nothing, because the rows it /// holds for a library like this name no file of ours (`record_in_place`). /// The bytes are in the library folder, and only the client may take them /// back. /// /// **This is a dehydration, not a deletion, and the distinction is the whole /// safety of the feature.** Removing a materialised file inside a synced tree /// propagates to the server and deletes the photograph everywhere. /// /// Best effort per image: a file the client refuses to release simply stays, /// which costs disk and loses nothing. pub fn spawn_dehydrate( conn: Connection, catalog_path: PathBuf, images: Vec, ) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "dehydrate", move || { let Ok(catalog) = Catalog::open(&catalog_path) else { return; }; let Ok(rt) = crate::net_runtime::build() else { return; }; rt.block_on(async { let Ok(backend) = crate::remote::connect(&conn) else { return; }; // Nothing to do where content is not a thing that can be given // back — a server library, or a plain folder. if !backend.capabilities().materialisation.can_materialise() { return; } let mut released = 0usize; for image in images { let Some(source_ref) = source_ref_of(&catalog, image) else { continue; }; let id = RemoteId::Path(RemotePath::new(&source_ref)); match backend.dematerialise(&id).await { Ok(()) => released += 1, Err(e) => log::debug!("releasing {source_ref}: {e}"), } } log::info!("released {released} photograph(s) back to the sync client"); let _ = tx.send(released); }); }); rx } /// What an image occupies, as the catalog recorded it. /// /// Zero where the scan could not tell — a placeholder reports no size, because /// a one-byte stub says nothing about what it stands for (ARCH §9.0a). A pin /// that cannot state its cost is better than one that states a wrong one. pub(super) fn size_of(catalog: &Catalog, image: dr_types::ImageId) -> Option { catalog .connection() .query_row( "SELECT file_size FROM images WHERE id = ?1", rusqlite::params![image.0 as i64], |r| r.get::<_, Option>(0), ) .ok() .flatten() .map(|v| v.max(0) as u64) } /// The remote path for a catalogued image. pub(super) fn source_ref_of(catalog: &Catalog, image: dr_types::ImageId) -> Option { catalog .connection() .query_row( "SELECT source_ref FROM images WHERE id = ?1", rusqlite::params![image.0 as i64], |r| r.get(0), ) .ok() } /// TRACES: FR-NC-6a | FR-CAT-9 /// Where a cached original is kept and how much may be kept. /// /// Passed in rather than derived here so the caller owns the policy: the /// budget is a user setting, and this function is on a worker thread with no /// access to one. pub struct CacheContext { pub dir: PathBuf, pub catalog_path: PathBuf, pub image: dr_types::ImageId, pub budget: dr_catalog::Budget, /// Whether a downloaded original is kept. /// /// Only the write. A cache is always *read*, because bytes already on disk /// cost nothing to use and declining them would re-download an image that /// is present — including every pinned one, which would leave a pinned /// collection unopenable offline the moment this was switched off. pub store: bool, } /// TRACES: FR-CAT-8 | FR-DEV-6 /// Fetch and parse the sidecar beside one image. /// /// # Why the edit is read from the file rather than the catalog /// /// The catalog carries a `graph_hash` and no parameters, and it is /// *disposable* (ARCH §6.12) — a rebuild would silently return every /// photograph to neutral. The sidecar is the authoritative store, so it is /// what an open reads, and that is also what makes an edit pasted on the /// desktop appear when the same frame is opened on the phone. /// /// # Why absence and failure are the same answer here /// /// `None` means "open this image at its defaults", which is right for a /// photograph that has never been edited — the overwhelmingly common case on a /// fresh library — and equally right when the network is down. The alternative, /// refusing to open the image because its sidecar could not be read, would make /// an unreachable server also mean an unviewable library. /// /// The one case that is *not* harmless is a sidecar that exists but does not /// parse. That still opens at defaults, but the write path /// ([`write_one_sidecar`]) independently refuses to overwrite a file it could /// not read, so an edit this build failed to understand is never destroyed by /// having been opened. pub fn spawn_sidecar_fetch( conn: Connection, image_path: String, cache_dir: PathBuf, offline: bool, ) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "sidecars", move || { let cache = SidecarCache::open(cache_dir); let path_str = sidecar_path(&image_path); // TRACES: FR-CAT-9 | FR-NC-10 // The cache wins outright when it is holding work the server has not // seen. Fetching in that state would answer with a document *older* // than the edit sitting in the outbox, and opening the photograph // would silently show it without the change the user just made — // which the next save would then write back over the top of. if cache.is_pending(&path_str) { log::debug!("{path_str} has queued local edits; opening from the cache"); let _ = tx.send(FetchedSidecar::unknown(cache.load(&path_str))); return; } // Offline there is nothing to ask, and the cache is the whole answer. let rt = if offline { None } else { match crate::net_runtime::build() { Ok(e) => Some(e), Err(e) => { log::debug!("sidecar fetch runtime: {e}"); None } } }; let Some(rt) = rt else { let _ = tx.send(FetchedSidecar::unknown(cache.load(&path_str))); return; }; rt.block_on(async { let backend = match crate::remote::connect(&conn) { Ok(b) => b, Err(e) => { log::debug!("sidecar fetch backend: {e}"); let _ = tx.send(FetchedSidecar::unknown(cache.load(&path_str))); return; } }; let path = RemotePath::new(path_str.clone()); let id = RemoteId::Path(path.clone()); // A 404 is the normal case on a library that has never been // edited, so this is `ok()` rather than an error path. let bytes = match backend.get(&id, None).await { Ok(bytes) => bytes, Err(e) => { // Unreachable, or no such file. The cache holds the best // answer this device has either way — but only a server // that said "no such file", with nothing in the cache, // is the word that there is no edit (FR-DEV-6). let cached = cache.load(&path_str); let absent = cached.is_none() && matches!(e, dr_sync::RemoteError::NotFound(_)); let _ = tx.send(FetchedSidecar { sidecar: cached, absent, }); return; } }; let text = String::from_utf8_lossy(&bytes).into_owned(); let parsed = match dr_pipeline::Sidecar::parse(&text) { Ok(mut s) => { // TRACES: FR-NC-8 | FR-NC-9 // Opening a photograph must show everything that has been // done to it, not whichever of two split default versions // happens to win. Fused in memory with no canonical uuid // to impose — this is a read, and the write path is where // the identity is decided. // // The fused document is what gets cached below, so the // next offline open sees the union too. s.fuse_default_versions(None); Some(s) } Err(e) => { log::warn!("sidecar at {} is unreadable ({e})", path.as_str()); None } }; // Populate the cache from what the server said, so the *next* // open of this photograph works with no connection. Clean rather // than pending: this content came from the server, so there is // nothing to send back. if let Some(sidecar) = parsed.as_ref() { if let Err(e) = cache.store(&path_str, sidecar, false) { log::debug!("caching {path_str}: {e}"); } } let _ = tx.send(FetchedSidecar::unknown(parsed)); }); }); rx } /// TRACES: FR-CAT-8 | FR-DEV-6 /// What a stored-edit fetch answered. /// /// `absent` is true only when the server said there is no such file and no /// cached copy stood in — the one answer that means "DarkRoom has no edit of /// this photograph", and so the only one that lets its earlier edit in. /// Every other way of arriving at no sidecar — offline, unreachable, a file /// that would not parse — leaves it false. #[derive(Debug, Default)] pub struct FetchedSidecar { pub sidecar: Option, pub absent: bool, } impl FetchedSidecar { fn unknown(sidecar: Option) -> Self { Self { sidecar, absent: false, } } } /// Fetch one file in full, for opening it in develop. /// /// Deliberately *not* the preview path. Browsing fetches a range and decodes /// an embedded JPEG (FR-NC-3); develop needs every byte, because demosaic /// needs every photosite. On a RAW file that is tens of megabytes, which is /// why this is a click-triggered download and not something the grid does. /// /// # Read-through /// /// With a `cache`, this checks disk before the network and stores what it /// downloads. That is what makes opening the same photograph twice cost one /// transfer, and what leaves a working session's images openable offline /// without anyone having pinned anything. /// /// A cache miss is not an error and a cache failure is not fatal: both fall /// through to the network, which is exactly the behaviour that existed before /// the cache did. /// /// Returns the bytes on a channel rather than blocking: the download runs on /// its own thread and the UI stays live, exactly as thumbnail fetching does. /// The work itself is [`fetch_original`], which is also what the /// [`Prefetcher`] runs — one transfer path, so a photograph fetched ahead is /// stored exactly as one fetched on a click. pub fn spawn_full_fetch( conn: Connection, path: String, cache: Option, ) -> Receiver, FetchFailure>> { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "fetch", move || { let _ = tx.send(fetch_original(conn, &path, cache.as_ref())); }); rx } /// TRACES: FR-NC-6a /// Every original transfer that is under way right now, by remote path. /// /// **One file, one transfer.** The [`Prefetcher`] fetches the photographs /// beside the open one before they are asked for, and the whole point is that /// the user then asks for one of them — often while it is still coming down. /// Without this the click would miss the cache, start a second download of /// the same file, and the two would halve each other's bandwidth for the rest /// of the transfer. With it, the click finds the path claimed, waits for the /// prefetch to store its bytes, and reads them from disk. /// /// Process-wide rather than passed in, because the property it enforces is /// process-wide: there is no caller for whom two concurrent downloads of one /// file is the right answer. Keyed on the path rather than the image id /// because that is the one name every caller has. static IN_FLIGHT: std::sync::LazyLock = std::sync::LazyLock::new(InFlight::default); /// TRACES: FR-NC-6a /// How far one original's transfer has got: bytes received, and the length /// the server declared (zero until it has, or if it never does). /// /// Held in the registry beside the claim rather than handed to the caller, /// because the caller watching is often not the one downloading: a step along /// the roll usually lands on a frame the [`Prefetcher`] is already fetching, /// and the click waits on that transfer instead of starting its own. #[derive(Default)] pub(super) struct Transfer { received: std::sync::atomic::AtomicU64, declared: std::sync::atomic::AtomicU64, } /// TRACES: FR-NC-6a /// Bytes received so far and bytes expected, for an original being fetched /// right now by anyone. `None` when nothing is fetching `path` — it is in the /// cache, or the transfer has not reached the network yet, or it has ended. pub fn transfer_progress(path: &str) -> Option<(u64, Option)> { IN_FLIGHT.progress(path) } #[derive(Default)] pub(super) struct InFlight { busy: std::sync::Mutex>>, freed: std::sync::Condvar, /// Threads parked in [`InFlight::claim`], counted under the lock so a /// test can release the holder only once a waiter is really waiting. #[cfg(test)] waiting: std::sync::atomic::AtomicUsize, } impl InFlight { /// Take `path` for this thread, or wait for whoever holds it. /// /// `Some` is a claim, released when the guard drops. `None` means another /// thread held the path and has now let it go — so the caller's cache /// check is worth repeating, because that thread has very probably just /// stored what the caller was about to download. fn claim(&self, path: &str) -> Option> { let mut busy = self.busy.lock().unwrap_or_else(|e| e.into_inner()); if !busy.contains_key(path) { let transfer = std::sync::Arc::new(Transfer::default()); busy.insert(path.to_string(), transfer.clone()); return Some(InFlightGuard { of: self, path: path.to_string(), transfer, }); } #[cfg(test)] self.waiting .fetch_add(1, std::sync::atomic::Ordering::SeqCst); while busy.contains_key(path) { busy = self.freed.wait(busy).unwrap_or_else(|e| e.into_inner()); } None } fn progress(&self, path: &str) -> Option<(u64, Option)> { use std::sync::atomic::Ordering::Relaxed; let busy = self.busy.lock().unwrap_or_else(|e| e.into_inner()); let t = busy.get(path)?; let declared = t.declared.load(Relaxed); Some((t.received.load(Relaxed), (declared > 0).then_some(declared))) } fn release(&self, path: &str) { self.busy .lock() .unwrap_or_else(|e| e.into_inner()) .remove(path); self.freed.notify_all(); } } /// A claim on a path, dropped on every exit from the fetch — a failed /// download must free the path too, or the click waiting on it never wakes. pub(super) struct InFlightGuard<'a> { of: &'a InFlight, path: String, transfer: std::sync::Arc, } impl Drop for InFlightGuard<'_> { fn drop(&mut self) { self.of.release(&self.path); } } /// The blocking body of [`spawn_full_fetch`]: cache, then in-flight registry, /// then network, storing what it downloads when the cache says to. pub(super) fn fetch_original( conn: Connection, path: &str, cache: Option<&CacheContext>, ) -> Result, FetchFailure> { // Opened on this thread: `rusqlite::Connection` is not `Send`, and the // UI thread's handle cannot be borrowed across the spawn. let cached = cache.and_then(|c| { let store = dr_catalog::Cache::open(&c.dir, c.budget).ok()?; let conn = Catalog::open(&c.catalog_path).ok()?; Some((store, conn)) }); let from_cache = || { let (c, (store, conn)) = (cache?, cached.as_ref()?); match store.load(conn.connection(), c.image, now_secs()) { Ok(Some(bytes)) => { log::info!("{path}: {} bytes from the local cache", bytes.len()); Some(bytes) } Ok(None) => None, // A cache that cannot be read is a cache miss, not a failure // to open the photograph. Err(e) => { log::debug!("cache lookup for {path}: {e}"); None } } }; // Miss, claim, and if the claim had to wait, look again: the thread that // held the path has finished with it, and what it fetched is on disk. let claim = loop { if let Some(bytes) = from_cache() { return Ok(bytes); } if let Some(claim) = IN_FLIGHT.claim(path) { break claim; } }; let rt = crate::net_runtime::build().map_err(FetchFailure::local)?; rt.block_on(async { let backend = crate::remote::connect(&conn).map_err(FetchFailure::local)?; let id = RemoteId::Path(RemotePath::new(path)); let transfer = &claim.transfer; let bytes = backend .get_reporting(&id, &|received, declared| { use std::sync::atomic::Ordering::Relaxed; transfer.received.store(received, Relaxed); transfer.declared.store(declared.unwrap_or(0), Relaxed); }) .await?; // Store before returning, so the bytes are on disk by the time the // image is on screen. Doing it after would leave a window where // closing the app immediately lost the download. if let (Some(c), Some((store, conn))) = (cache.filter(|c| c.store), cached.as_ref()) { // `pinned: false` — this is the passive population. A pin is // something the user asks for explicitly; opening an image is // not that, and treating it as one would make the pinned set // grow silently and never be evicted. if let Err(e) = store.store(conn.connection(), c.image, path, &bytes, false, now_secs()) { log::debug!("caching {path}: {e}"); } else if let Err(e) = store.enforce(conn.connection()) { log::debug!("enforcing the cache budget: {e}"); } } Ok(bytes) }) } /// TRACES: FR-NC-6a | FR-UI-4 /// One original to fetch ahead of its being asked for. pub struct PrefetchJob { pub path: String, pub cache: CacheContext, } /// What the [`Prefetcher`] is doing, for the activity list. pub enum PrefetchEvent { /// A transfer has started for this path. Started(String), /// And has ended — stored, or not; either way the row can go. Ended(String), } /// TRACES: FR-NC-6a | FR-UI-4 /// Fetches the photographs beside the open one into the cache, ahead of the /// step that asks for them. /// /// **Why this exists.** Walking the photo roll is one click per frame, and /// without this every click is a download of tens of megabytes with a /// "Downloading…" line over an empty canvas. A photographer moving between a /// pair of near-identical frames does that a dozen times. Fetching the two /// neighbours while the current photograph is being looked at turns the next /// step into a disk read, which is what makes stepping feel like stepping. /// /// **One worker, one wish.** A single thread serves the *latest* request and /// nothing older. Each open replaces the previous wish outright, so a fast /// walk along the roll does not leave a trail of stale downloads competing /// with the one the user is actually waiting on; a job already under way is /// finished rather than abandoned, because the bytes are mostly here. Jobs /// run in the order given — next before previous, since that is the way a /// roll is mostly walked — and one at a time, so two neighbours never halve /// each other's bandwidth. /// /// **What it never does.** It never fetches into a cache that would not keep /// the bytes: the caller only hands it jobs whose cache stores, because a /// prefetch that is discarded on arrival is pure transfer for nothing — and /// "keep opened originals" being off is the user saying this device is /// metered or small (FR-NC-6). It never starts while offline, for the same /// reason. And it fetches only the immediate neighbours: originals are /// "explicit pin or on-demand open only", and ±1 is as far as "on demand" /// honestly stretches. pub struct Prefetcher { shared: std::sync::Arc, events: Receiver, } #[derive(Default)] pub(super) struct PrefetchShared { wanted: std::sync::Mutex, changed: std::sync::Condvar, } /// The latest wish, and a generation so the worker can tell it has been /// replaced mid-list. #[derive(Default)] pub(super) struct Wanted { generation: u64, conn: Option, jobs: Vec, } impl Wanted { /// Replace whatever was wanted with `jobs`. fn replace(&mut self, conn: Connection, jobs: Vec) { self.generation += 1; self.conn = Some(conn); self.jobs = jobs; } /// Whether a wish taken at `generation` is still the current one. fn is_current(&self, generation: u64) -> bool { self.generation == generation } } impl Default for Prefetcher { fn default() -> Self { Self::new() } } impl Prefetcher { /// Start the worker. It sleeps until the first [`Self::want`]. pub fn new() -> Self { let shared = std::sync::Arc::new(PrefetchShared::default()); let (tx, events) = std::sync::mpsc::channel(); let worker = shared.clone(); executors::try_spawn(Executor::Network, "prefetch", move || { serve_prefetches(&worker, &tx) }) .expect("spawning the prefetch worker"); Self { shared, events } } /// Fetch these, in this order, instead of whatever was asked for before. /// /// An empty list is a valid wish: it cancels the rest of the previous /// one, and is what a photograph with no neighbours in the window asks. pub fn want(&self, conn: Connection, jobs: Vec) { self.shared .wanted .lock() .unwrap_or_else(|e| e.into_inner()) .replace(conn, jobs); self.shared.changed.notify_one(); } /// Everything the worker has reported since the last poll. pub fn poll(&self) -> Vec { std::iter::from_fn(|| self.events.try_recv().ok()).collect() } } /// The worker: take the current wish, serve it job by job, stop the moment it /// is superseded, sleep until the next one. pub(super) fn serve_prefetches(shared: &PrefetchShared, events: &Sender) { loop { let (generation, conn, jobs) = { let mut wanted = shared.wanted.lock().unwrap_or_else(|e| e.into_inner()); while wanted.jobs.is_empty() { wanted = shared .changed .wait(wanted) .unwrap_or_else(|e| e.into_inner()); } let jobs = std::mem::take(&mut wanted.jobs); let Some(conn) = wanted.conn.clone() else { continue; }; (wanted.generation, conn, jobs) }; // One catalog connection for this wish's row checks, opened at the // first and kept for the rest, rather than one per neighbour. let mut catalog = None; for job in jobs { let current = shared .wanted .lock() .unwrap_or_else(|e| e.into_inner()) .is_current(generation); if !current { break; } if holds_original(&job.cache, &mut catalog) { continue; } // A closed channel means the window is gone: nothing to fetch // for any more. if events .send(PrefetchEvent::Started(job.path.clone())) .is_err() { return; } match fetch_original(conn.clone(), &job.path, Some(&job.cache)) { Ok(bytes) => log::info!("{}: {} bytes fetched ahead", job.path, bytes.len()), // Not a failure anyone needs to hear about now: the click // that wants this photograph will try again and say so. Err(e) => log::debug!("fetching {} ahead: {}", job.path, e.message), } if events.send(PrefetchEvent::Ended(job.path)).is_err() { return; } } } } /// Whether the cache already has this original — a row check, not a read, /// so asking costs nothing and touches no `last_used`. /// /// Asked on `held`, which the caller keeps across a batch of these: a /// connection per question was an open per prefetched neighbour, four of the /// five a develop landing made. Opened here on first use, and again only if /// the batch moves to another catalog. pub(super) fn holds_original(cache: &CacheContext, held: &mut Option<(PathBuf, Catalog)>) -> bool { let Ok(store) = dr_catalog::Cache::open(&cache.dir, cache.budget) else { return false; }; if held .as_ref() .is_none_or(|(path, _)| *path != cache.catalog_path) { *held = Catalog::open(&cache.catalog_path) .ok() .map(|c| (cache.catalog_path.clone(), c)); } let Some((_, catalog)) = held.as_ref() else { return false; }; store.holds_original(catalog.connection(), cache.image) } #[cfg(test)] mod tests { use super::*; /// The whole reason the registry exists: a click on a photograph that is /// being fetched ahead waits for that transfer rather than starting its /// own, and is told so — `None` — so it looks in the cache again. #[test] fn a_second_claim_waits_for_the_first_to_be_released() { let registry = std::sync::Arc::new(InFlight::default()); let first = registry.claim("shoot/one.CR2"); assert!(first.is_some(), "an unclaimed path is claimed outright"); let waiter = { let registry = registry.clone(); std::thread::spawn(move || registry.claim("shoot/one.CR2").is_some()) }; // Wait until the second claim is parked on the condvar. The count is // bumped under the lock just before the wait, and the release below // needs that lock, so it cannot reach the waiter any earlier. while registry.waiting.load(std::sync::atomic::Ordering::SeqCst) == 0 { std::thread::yield_now(); } assert!( !waiter.is_finished(), "the second claim must not return while the first is held" ); drop(first); let claimed = waiter.join().unwrap(); assert!( !claimed, "after waiting, the caller is told to recheck the cache" ); assert!( registry.claim("shoot/one.CR2").is_some(), "and once nobody holds the path it can be claimed again" ); } /// Whoever is watching a path reads the holder's progress through the /// registry, and loses it when the transfer ends — a step onto a frame /// the prefetcher is fetching shows that transfer's bar. #[test] fn progress_is_readable_by_path_while_claimed() { use std::sync::atomic::Ordering::Relaxed; let registry = InFlight::default(); assert_eq!(registry.progress("a.CR2"), None); let claim = registry.claim("a.CR2").unwrap(); assert_eq!(registry.progress("a.CR2"), Some((0, None))); claim.transfer.received.store(1024, Relaxed); claim.transfer.declared.store(4096, Relaxed); assert_eq!(registry.progress("a.CR2"), Some((1024, Some(4096)))); assert_eq!(registry.progress("b.CR2"), None); drop(claim); assert_eq!(registry.progress("a.CR2"), None); } /// Different photographs never wait on each other. #[test] fn distinct_paths_are_claimed_independently() { let registry = InFlight::default(); let _a = registry.claim("a.CR2"); assert!(registry.claim("b.CR2").is_some()); } /// A newer wish supersedes an older one mid-list; the worker checks this /// between jobs, and it is what keeps a fast walk along the roll from /// queueing every neighbour it passed. #[test] fn a_new_wish_supersedes_the_one_being_served() { let mut wanted = Wanted::default(); let conn = test_connection(); wanted.replace(conn.clone(), Vec::new()); let taken = wanted.generation; assert!(wanted.is_current(taken)); wanted.replace(conn, Vec::new()); assert!( !wanted.is_current(taken), "the list taken before the replacement is stale" ); } fn test_connection() -> Connection { Connection::new( Account::new("nextcloud", "https://cloud.example").with_login("d", "d"), None, ) } }