Files
DarkRoom/ui/dr-ui/src/library/thumbnails_fetch.rs
T
dtourolle b1d1c47261 Start every worker thread through the executors module
Thirty-nine spawn sites in dr-ui, and one in the Android entry point,
called std::thread::spawn or a Builder of their own, and most of the
threads they started were <unnamed> in a panic message or a profiler.
Each now calls executors::spawn with its executor and a role, so the thread is
named <executor>:<role> — net:sync, decode:thumbs, io:catalog-open —
and knows which executor it is on. The three that already set a name
(automation, import, prefetch) keep their name as the role.

Behaviour is unchanged: each job still gets a thread of its own when it
starts, and spawn panics where std::thread::spawn did.

The module's documentation now says how a job is assigned: by what it
spends its time on, so a sweep that fetches bytes and then decodes them
is Decode, and a sidecar write that touches the catalog is Network.

Left as they were: the segmentation and refine workers in masks_ui.rs,
which another change is reworking, and test-only threads.
2026-09-27 07:08:37 -04:00

1051 lines
39 KiB
Rust

//! 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<u64>,
/// 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<dr_sync::RemoteError> 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<PinMessage> {
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<dr_types::ImageId>,
) -> Receiver<usize> {
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<u64> {
catalog
.connection()
.query_row(
"SELECT file_size FROM images WHERE id = ?1",
rusqlite::params![image.0 as i64],
|r| r.get::<_, Option<i64>>(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<String> {
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<Option<dr_pipeline::Sidecar>> {
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(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(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(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 Ok(bytes) = backend.get(&id, None).await else {
// Unreachable, or no such file. The cache cannot tell those
// apart and does not need to: either way it holds the best
// answer this device has.
let _ = tx.send(cache.load(&path_str));
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(parsed);
});
});
rx
}
/// 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<CacheContext>,
) -> Receiver<Result<Vec<u8>, 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<InFlight> = 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<u64>)> {
IN_FLIGHT.progress(path)
}
#[derive(Default)]
pub(super) struct InFlight {
busy: std::sync::Mutex<std::collections::HashMap<String, std::sync::Arc<Transfer>>>,
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<InFlightGuard<'_>> {
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<u64>)> {
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<Transfer>,
}
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<Vec<u8>, 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<PrefetchShared>,
events: Receiver<PrefetchEvent>,
}
#[derive(Default)]
pub(super) struct PrefetchShared {
wanted: std::sync::Mutex<Wanted>,
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<Connection>,
jobs: Vec<PrefetchJob>,
}
impl Wanted {
/// Replace whatever was wanted with `jobs`.
fn replace(&mut self, conn: Connection, jobs: Vec<PrefetchJob>) {
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<PrefetchJob>) {
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<PrefetchEvent> {
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<PrefetchEvent>) {
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,
)
}
}