Files
DarkRoom/ui/dr-ui/src/library/thumbnails_fetch.rs
T
dtourolle a1d511fd4b Split library.rs into library/ by area of behaviour
library.rs was 7,729 lines wiring together everything "open a remote
library" touches: scanning, pulling other devices' judgements out of
sidecars found along the way, writing local edits back out to the
sidecar outbox, pushing/reloading XMP by hand, fetching and prefetching
thumbnails and originals, generating thumbnails locally, the metadata
and thumbnail background sweeps, on-disk paths for the catalog and
model files, and reading the grid's cells, spans and rating filter.
Same motivation as the develop.rs split (docs/dev/code-health.md CH-1):
a pure, no-behaviour-change move into one file per area, each under
about 1,500 lines.

Tracing actual call sites rather than trusting the file's physical
layout mattered here: `persist`, `load_folder_etags`, `pull_sidecars`,
`load_sidecar_etags`, `record_sidecar_read` and `apply_judgement` sit
textually beside the XMP push/reload functions but are called only
from `run_scan` (pulling a device's own past judgements out of the
sidecars a scan just walked), so they went to scan.rs and not xmp.rs.
`cells` came out at over 1,800 lines once its tests moved with it and
split further into cells.rs (windowed reads, trash, ordinals) and
spans.rs (collection scope, manual reordering, the capture-time
histogram) -- ten submodules rather than the nine first planned.

Previously-private items reached from a sibling module became
`pub(super)`, narrower than the whole-crate reachability one file gave
them. Tests moved with the code they test; the two test fixtures used
across more than one file (`scanned`, and develop.rs's
`session_with_a_left_half_subject` in the matching commit) joined the
shared `test_support` module alongside the existing `entry`/
`with_images`/`image_ids` helpers. `mod.rs` re-exports every module's
public items under `library::`, including the `pub(crate)`
`test_support` module `repairs.rs` reads its fixtures from, so no file
outside `library` needed a change.

The previous commit split develop.rs the same way; taken alone it left
dr-ui without library.rs, so that intermediate commit does not build on
its own. This one restores it.
2026-09-20 18:21:43 +02:00

966 lines
36 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::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();
std::thread::spawn(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();
std::thread::spawn(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();
std::thread::spawn(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();
std::thread::spawn(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);
#[derive(Default)]
pub(super) 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.
pub(super) 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.
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 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}");
}
}
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();
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.
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)
};
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;
}
// 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`.
pub(super) 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)
}
#[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 (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,
)
}
}