Fetch the photographs around the open one ahead of the step to them
Benchmarks / CPU and I/O (per commit) (push) Successful in 3m59s
Benchmarks / Frame budget (on demand) (push) Skipped
Build and test / android-image (push) Canceled after 0s
🐳 Android image / Build and push (push) Canceled after 0s
Build and test / Android (aarch64) (push) Canceled after 0s
Build and test / windows-image (push) Canceled after 0s
🐳 Windows image / Build and push (push) Canceled after 0s
Build and test / Windows (x86_64, cross) (push) Canceled after 0s
Build and test / Layer separation (push) Canceled after 0s
Build and test / Desktop (Linux) (push) Canceled after 16m26s
Traceability / Requirement traces (push) Canceled after 0s

Walking the photo roll was one download per frame: every step showed
"Downloading…" over an empty canvas while tens of megabytes came down,
and moving between a pair of near-identical frames paid that a dozen
times. Now, once the opened photograph has landed, the ones around it
are fetched into the originals cache while it is being looked at, so
the next step is a disk read.

A single worker serves the latest wish only, closest first and working
outwards — next, previous, next-but-one, previous-but-one… — one file
at a time. Each open replaces the wish, so a fast walk never leaves a
trail of stale downloads competing with the one being waited on. A
process-wide in-flight registry makes a click on a photograph that is
still being fetched ahead wait for that transfer and read it from disk,
rather than start a second download of the same file.

How far each side is a setting under STORAGE — Off, 2, 5, 10 or 20,
defaulting to 5 — and it is moot while "keep originals after opening"
is off, since a fetch the cache would discard on arrival is transfer
for nothing. Nothing is fetched ahead while offline. The transfers show
in the activity list while they run and are removed when they end.
This commit is contained in:
2026-09-19 10:36:50 +02:00
parent 2917b7427d
commit 78cb00634e
9 changed files with 775 additions and 154 deletions
+409 -77
View File
@@ -2424,26 +2424,6 @@ pub struct CacheContext {
pub store: bool,
}
/// 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.
/// TRACES: FR-CAT-8 | FR-DEV-6
/// Fetch and parse the sidecar beside one image.
///
@@ -2570,80 +2550,356 @@ pub fn spawn_sidecar_fetch(
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();
std::thread::spawn(move || {
// Opened on this thread: `rusqlite::Connection` is not `Send`, and the
// UI thread's handle cannot be borrowed across the spawn.
let cached = cache.as_ref().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 _ = tx.send(fetch_original(conn, &path, cache.as_ref()));
});
rx
}
if let (Some(c), Some((store, conn))) = (cache.as_ref(), 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());
let _ = tx.send(Ok(bytes));
return;
}
Ok(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}"),
/// 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);
#[derive(Default)]
struct InFlight {
busy: std::sync::Mutex<std::collections::HashSet<String>>,
freed: std::sync::Condvar,
}
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.insert(path.to_string()) {
return Some(InFlightGuard {
of: self,
path: path.to_string(),
});
}
while busy.contains(path) {
busy = self.freed.wait(busy).unwrap_or_else(|e| e.into_inner());
}
None
}
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.
struct InFlightGuard<'a> {
of: &'a InFlight,
path: String,
}
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.
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 bytes = backend.get(&id, None).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}");
}
}
let rt = match crate::net_runtime::build() {
Ok(rt) => rt,
Err(e) => {
let _ = tx.send(Err(FetchFailure::local(e)));
return;
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)]
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)]
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();
std::thread::Builder::new()
.name("prefetch".into())
.spawn(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.
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)
};
rt.block_on(async {
let backend = match crate::remote::connect(&conn) {
Ok(b) => b,
Err(e) => {
let _ = tx.send(Err(FetchFailure::local(e)));
return;
}
};
let id = RemoteId::Path(RemotePath::new(&path));
let got = backend.get(&id, None).await.map_err(FetchFailure::from);
// Store before sending, 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 (Ok(bytes), Some(c), Some((store, conn))) =
(&got, cache.as_ref().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}");
}
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) {
continue;
}
let _ = tx.send(got);
});
});
// 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;
}
}
}
}
rx
/// Whether the cache already has this original — a row check, not a read,
/// so asking costs nothing and touches no `last_used`.
fn holds_original(cache: &CacheContext) -> bool {
let Ok(store) = dr_catalog::Cache::open(&cache.dir, cache.budget) else {
return false;
};
let Ok(catalog) = Catalog::open(&cache.catalog_path) else {
return false;
};
store.holds_original(catalog.connection(), cache.image)
}
/// Serve thumbnails for a set of rows: store first, network second.
@@ -7864,3 +8120,79 @@ mod amending_across_devices {
assert_eq!(out.versions["for-print"].rating, 5);
}
}
#[cfg(test)]
mod fetching_ahead {
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 (tx, rx) = std::sync::mpsc::channel();
let waiter = {
let registry = registry.clone();
std::thread::spawn(move || {
tx.send(()).unwrap();
registry.claim("shoot/one.CR2").is_some()
})
};
rx.recv().unwrap();
// The waiter is blocked on the first claim. Not provable without a
// sleep, but a release that reaches it proves the wait ended there.
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"
);
}
/// 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,
)
}
}