Files
DarkRoom/ui/dr-ui/src/library/sweep.rs
T
dtourolle 64ea44aefe Stop reading dates at the first sign the server is unreachable
A window of cells whose thumbnails were cached but whose dates were not
sent a header read per cell, and offline each one was three attempts at
a 15 s connect timeout: `is_transient` counts a network error as worth
retrying, and `read_metadata_only` returned a bare bool that could not
say why a read failed. So the grid sat on "reading N dates" for minutes
against a server that was not there, and no banner went up, because
nothing in that loop ever reported the connection.

`read_metadata_only` now returns a `DateRead`: reached, failed, or
offline. An offline error is returned on the first attempt rather than
retried — a dead server answers the second exactly as the first — while
a 423 lock is still retried, which is what the retry was for. The grid's
worker stops on it and sends `Offline`, as its fetch loop already did,
so the banner goes up and the bar stops. The sweep's lanes stop on it
too, one timeout each rather than one per image.
2026-10-03 16:49:49 -04:00

1130 lines
45 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! The background sweeps: metadata extraction and thumbnail generation
//! for whatever the catalog still owes, a chunk at a time.
use crate::executors::{self, Executor};
use dr_catalog::Catalog;
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
use dr_thumbs::ThumbStore;
use std::path::PathBuf;
use std::sync::mpsc::{Receiver, Sender};
use super::filters::{VISIBLE, VISIBLE_UNALIASED};
use super::thumbnails_fetch::ThumbnailRequest;
use super::thumbnails_gen::{
collect_metadata, encode_preview, fetch_preview, store_thumbnail, PreviewOutcome,
ThumbnailMessage,
};
/// Capture metadata read from the same header the thumbnail needed.
///
/// Free: the header fetch happens either way, so parsing EXIF out of it costs
/// no extra transfer. That is what fills the timeline as the user browses,
/// rather than a separate 6 GB sweep over the library.
#[derive(Debug, Clone)]
pub struct MetadataFound {
pub image_id: i64,
pub captured_at: Option<i64>,
pub captured_offset: Option<i32>,
pub camera: Option<String>,
pub lens: Option<String>,
pub iso: Option<u32>,
/// TRACES: FR-MRG-6
/// Width and height as the photograph is seen — orientation applied —
/// where the header says: what gives a panorama from elsewhere, a
/// stitch from another program or a phone's sweep, its wide cell.
pub size: Option<(u32, u32)>,
}
impl MetadataFound {
/// The upright size a header describes.
pub fn upright_size(md: &dr_decode::Metadata) -> Option<(u32, u32)> {
let (w, h) = (md.width?, md.height?);
if w == 0 || h == 0 {
return None;
}
let turned = md.orientation.is_some_and(|o| o.quarter_turns % 2 == 1);
Some(if turned { (h, w) } else { (w, h) })
}
}
/// Write a batch of dates and tell the UI, draining `found`.
///
/// Separate from the loop so the same path serves both the periodic flush and
/// the final one, and so a write failure is reported once rather than being
/// silently swallowed by the caller.
pub(super) fn flush_metadata(
catalog_path: &std::path::Path,
found: &mut Vec<MetadataFound>,
tx: &Sender<ThumbnailMessage>,
) {
if found.is_empty() {
return;
}
match Catalog::open(catalog_path) {
Ok(cat) => match write_metadata(&cat, found) {
Ok(n) => {
log::info!("recorded capture dates for {n} of {} image(s)", found.len());
// Tell the UI so the timeline can appear. Without this the
// histogram only shows up on the next window load, which on a
// fully cached library may be never.
let _ = tx.send(ThumbnailMessage::DatesRecorded(n));
}
Err(e) => log::warn!("writing metadata: {e}"),
},
Err(e) => log::warn!("opening catalog to write metadata: {e}"),
}
found.clear();
}
/// Read only the date for an image whose thumbnail is already cached.
///
/// One 256 KB header request, no preview range and no decode. This is what
/// gets a library dated when its thumbnails came from the store — including
/// shards synced from another device, which carry pixels but no metadata.
/// Read a header for its date.
///
/// Reports whether the file was **reached**, which the caller needs and cannot
/// otherwise tell: a header that carried no EXIF and a fetch that never
/// happened both leave `found` untouched, and recording the second as "this
/// image has no date" would let one lock mark it dateless for good.
pub(crate) async fn read_metadata_only(
backend: &dyn RemoteBackend,
decoder: &dyn dr_decode::Decoder,
req: &ThumbnailRequest,
found: &mut Vec<MetadataFound>,
) -> DateRead {
let id = RemoteId::Path(RemotePath::new(&req.path));
// Retried, because one failure here is usually a lock rather than a
// verdict. Nextcloud's file locking answers a plain *read* with 423 under
// concurrency, and the identical range succeeds moments later — measured
// against a real server while twelve lanes were running. Without a retry
// those images sit out the whole pass over a lock that lasted a moment.
//
// Bounded and short: a genuinely missing or forbidden file must not cost
// three round trips before the sweep moves on.
const ATTEMPTS: usize = 3;
for attempt in 1..=ATTEMPTS {
match backend.get(&id, Some(0..decoder.header_bytes())).await {
Ok(header) => {
collect_metadata(backend, decoder, &id, &header, req, found).await;
return DateRead::Reached;
}
// Not retried, unlike a lock. A dead server answers the second
// attempt exactly as it answered the first, after the same 15 s
// connect timeout — three of those per image turned a window of
// cached-but-undated cells into minutes of "reading dates"
// against nothing, with no banner, because nothing said why.
Err(e) if e.indicates_offline() => {
log::debug!("reading date for {}: {e}", req.path);
return DateRead::Offline(e.to_string());
}
Err(e) if e.is_transient() && attempt < ATTEMPTS => {
// Backing off at all matters more than the exact interval: the
// contention that produced the lock is our own lanes, so any
// pause lets the holder finish.
tokio::time::sleep(std::time::Duration::from_millis(200 * attempt as u64)).await;
}
Err(e) => {
// Not surfaced: a missing date leaves the image off the
// timeline rather than breaking anything, and the next sweep
// retries it regardless.
log::debug!("reading date for {} ({attempt} attempts): {e}", req.path);
return DateRead::Failed;
}
}
}
DateRead::Failed
}
/// What one header read for a date came to.
#[derive(Debug, PartialEq, Eq)]
pub(crate) enum DateRead {
/// The header arrived. Whatever EXIF it held is in `found`.
Reached,
/// This file could not be read; the next pass tries it again.
Failed,
/// The server could not be reached, so no read after this one will be
/// either. The caller stops rather than paying a timeout per image.
Offline(String),
}
/// Write capture metadata read during the thumbnail pass.
///
/// Promotes each row from `metadata_state = 1` (stat-only) to 2 (full EXIF),
/// which is what makes it eligible for the timeline. A row whose EXIF was
/// unreadable stays at 1 rather than being marked done with empty fields, so a
/// later attempt can retry it.
///
/// Returns how many rows were promoted.
pub fn write_metadata(
catalog: &Catalog,
found: &[MetadataFound],
) -> Result<usize, dr_catalog::CatalogError> {
let conn = catalog.connection();
let tx = conn.unchecked_transaction()?;
let mut promoted = 0;
for m in found {
// Only a real timestamp counts as fully read. Camera and lens without
// a date leave the image unplaceable on a timeline, which is exactly
// the state the grid needs to distinguish.
let state = if m.captured_at.is_some() { 2 } else { 1 };
tx.execute(
"UPDATE images
SET captured_at = coalesce(?2, captured_at),
captured_offset = coalesce(?3, captured_offset),
camera = coalesce(?4, camera),
lens = coalesce(?5, lens),
iso = coalesce(?6, iso),
metadata_state = max(metadata_state, ?7),
w = coalesce(?8, w),
h = coalesce(?9, h)
WHERE id = ?1",
rusqlite::params![
m.image_id,
m.captured_at,
m.captured_offset,
m.camera,
m.lens,
m.iso,
state,
m.size.map(|s| s.0),
m.size.map(|s| s.1),
],
)?;
if state == 2 {
promoted += 1;
}
}
tx.commit()?;
Ok(promoted)
}
/// Progress from the whole-library sweep.
#[derive(Debug)]
pub enum SweepMessage {
/// How many images still need work, counted once at the start.
Total(usize),
/// Another chunk finished. Carries cumulative counts.
Progress {
done: usize,
dated: usize,
},
Finished {
dated: usize,
},
}
/// Await every future concurrently, returning results in order.
///
/// A hand-rolled `join_all` rather than a `futures` dependency for one
/// function. Polling a `Vec` of futures in a loop is exactly what the crate's
/// version does; the ordering guarantee is what lets the caller pair results
/// back to their inputs.
pub(crate) async fn futures_join_all<F>(futures: impl IntoIterator<Item = F>) -> Vec<F::Output>
where
F: std::future::Future,
{
use std::pin::Pin;
use std::task::Poll;
// Boxed so each future has a stable address while it is polled in place.
let mut pending: Vec<Option<Pin<Box<F>>>> =
futures.into_iter().map(|f| Some(Box::pin(f))).collect();
let mut done: Vec<Option<F::Output>> = (0..pending.len()).map(|_| None).collect();
std::future::poll_fn(move |cx| {
let mut all_ready = true;
for (slot, out) in pending.iter_mut().zip(done.iter_mut()) {
let Some(fut) = slot else { continue };
match fut.as_mut().poll(cx) {
Poll::Ready(v) => {
*out = Some(v);
// Dropped as soon as it completes, so a long-running lane
// does not hold a finished one's resources.
*slot = None;
}
Poll::Pending => all_ready = false,
}
}
if all_ready {
Poll::Ready(done.iter_mut().filter_map(Option::take).collect())
} else {
Poll::Pending
}
})
.await
}
/// How many images one sweep chunk handles before committing.
///
/// Small enough that a kill loses little, large enough that the catalog is not
/// reopened per image. A multiple of [`SWEEP_LANES`] so every lane gets equal
/// work and no chunk ends with most lanes idle.
pub(super) const SWEEP_CHUNK: usize = 96;
/// How many fetches the sweep keeps in flight.
///
/// Each is ~0.6 s of round-trip latency and almost no bandwidth — a 256 KB
/// header — so the sequential version spent essentially all its time waiting.
/// Twelve lanes turn ~3 hours into ~15 minutes on the reference library.
///
/// Deliberately bounded rather than unlimited: the grid's own interactive
/// fetches share this server, and a sweep that saturated the connection would
/// make browsing feel broken while it ran.
///
/// **Lowered from twelve after measuring.** Twelve produced 423 Locked on a
/// real server — Nextcloud's file locking answering a plain read under
/// contention we were creating ourselves. Six keeps most of the speedup
/// without provoking it; the retry above covers what still slips through.
pub(crate) const SWEEP_LANES: usize = 6;
/// TRACES: FR-CULL-8 | NFR-RES-2
/// The largest original the face sweep will fetch, in bytes.
///
/// A budget, not a correctness rule: nothing about a byte count says whether
/// a file decodes. It exists because the sweep fetches the whole original
/// before it can learn anything about it, and the one file in the reference
/// library above this line is a 521 MB stitched panorama the decoder refuses
/// on sight — so every pass on the tablet spent half a gigabyte of Wi-Fi to
/// find that out again. Below the line: every camera RAW this library holds,
/// the largest a 60 MB medium-format file; above it, four files, all
/// panoramas.
///
/// A file over budget is marked examined with nothing found and a zero
/// edge, the same mark a file the decoder cannot open gets, so the count in
/// the sweep's report says it was skipped and a later pass can select it.
/// That later pass is the real answer for a panorama — read it in tiles,
/// detect in each, and stitch the boxes back — and this constant is the
/// placeholder for it, not a decision that panoramas hold no faces.
pub(crate) const SWEEP_MAX_ORIGINAL_BYTES: u64 = 256 * 1024 * 1024;
/// Date **every** image in the library, not just the ones on screen.
///
/// Thumbnails are deliberately *not* fetched here. A thumbnail needs the
/// mutable store, which cannot be shared across the parallel lanes below, and
/// it costs 1–3 MB against a date's 256 KB. Dating the whole library is what
/// the timeline needs; thumbnails arrive as cells are actually browsed, which
/// is the FR-NC-3 posture anyway.
///
/// The grid's own fetches cover what is on screen; this covers the rest, so the
/// timeline describes the whole library rather than the part that happened to
/// be scrolled past. It is resumable by construction — each pass queries for
/// what is still missing, so a kill mid-sweep costs only the current chunk.
///
/// Runs at the back of the queue by design: it holds no lock the grid needs,
/// and its chunked commits keep write transactions short.
pub fn spawn_sweep(conn: Connection, catalog_path: PathBuf) -> Receiver<SweepMessage> {
let (tx, rx) = std::sync::mpsc::channel();
// TRACES: FR-RAW-2
// The one place this job names a decoder; everything below takes it.
let decoder = dr_decode::default();
executors::spawn(Executor::Decode, "metadata", move || {
let catalog = match Catalog::open(&catalog_path) {
Ok(c) => c,
Err(e) => {
// Silent failure here left the sweep looking like it had run
// and found nothing: no progress, no error, 17,397 images
// still unindexed.
log::warn!(
"sweep: cannot open catalog at {}: {e}",
catalog_path.display()
);
let _ = tx.send(SweepMessage::Finished { dated: 0 });
return;
}
};
let outstanding = count_outstanding(&catalog).unwrap_or(0);
if outstanding == 0 {
let _ = tx.send(SweepMessage::Finished { dated: 0 });
return;
}
log::info!("sweep: {outstanding} image(s) need a date or a thumbnail");
if tx.send(SweepMessage::Total(outstanding)).is_err() {
return;
}
let rt = match crate::net_runtime::build() {
Ok(rt) => rt,
Err(e) => {
log::warn!("sweep: no runtime: {e}");
return;
}
};
rt.block_on(async {
let Ok(backend) = crate::remote::connect(&conn) else {
return;
};
let (mut done, mut dated) = (0usize, 0usize);
loop {
// Re-queried each pass rather than held as one long list: the
// grid is dating images at the same time, and a stale list
// would refetch what it already covered.
let chunk = match next_outstanding(&catalog, SWEEP_CHUNK) {
Ok(c) if c.is_empty() => break,
Ok(c) => c,
Err(e) => {
log::warn!("sweep: {e}");
break;
}
};
let chunk_started = std::time::Instant::now();
log::debug!(
"sweep: chunk of {} starting at image {}",
chunk.len(),
chunk[0].image_id
);
// Twelve lanes over the chunk. Each lane owns a disjoint slice
// and its own `found` vector, so nothing is shared and no lock
// is needed; the results are concatenated after the join.
//
// The thumbnail store is the exception — it is `&mut` and
// cannot be shared — so lanes only *read* metadata and any
// missing thumbnail is left to the interactive path. Dating the
// library is what the sweep is for; thumbnails arrive as cells
// are browsed.
let lanes: Vec<Vec<&ThumbnailRequest>> = (0..SWEEP_LANES)
.map(|lane| chunk.iter().skip(lane).step_by(SWEEP_LANES).collect())
.collect();
let results = futures_join_all(lanes.into_iter().map(|lane| {
let backend: &dyn RemoteBackend = &*backend;
async move {
let mut found = Vec::new();
let mut reached = Vec::new();
for req in lane {
match read_metadata_only(backend, decoder, req, &mut found).await {
DateRead::Reached => reached.push(req.image_id),
DateRead::Failed => {}
// The other lanes find the same, each after
// one timeout rather than one per image.
DateRead::Offline(_) => break,
}
}
(found, reached)
}
}))
.await;
let mut found = Vec::new();
let mut reached = std::collections::HashSet::new();
for (lane_found, lane_reached) in results {
found.extend(lane_found);
reached.extend(lane_reached);
}
done += chunk.len();
// An image whose header carried no EXIF at all yields nothing
// to `found`, so nothing marks it examined and the next sweep
// fetches it again — for ever. Darktable exports strip
// metadata by default, and 2,188 of them in the reference
// library meant 2,188 pointless round trips per run.
//
// Recorded as examined with no date: the file was read and
// genuinely has none, which is a different state from "not
// looked at yet" and must not be confused with it.
let answered: std::collections::HashSet<i64> =
found.iter().map(|m| m.image_id).collect();
// Only files actually read. One that could not be fetched is
// left alone so the next pass retries it, rather than being
// written off over a lock or a dropped connection.
found.extend(
chunk
.iter()
.filter(|r| {
reached.contains(&r.image_id) && !answered.contains(&r.image_id)
})
.map(|r| MetadataFound {
image_id: r.image_id,
captured_at: None,
captured_offset: None,
camera: None,
lens: None,
iso: None,
size: None,
}),
);
dated += found.iter().filter(|m| m.captured_at.is_some()).count();
let read = answered.len();
flush_sweep(&catalog, &mut found);
log::info!(
"sweep: {done} done, {dated} dated ({read} read in {:.1}s)",
chunk_started.elapsed().as_secs_f64()
);
if tx.send(SweepMessage::Progress { done, dated }).is_err() {
return;
}
}
log::info!("sweep complete: {dated} date(s) recorded over {done} image(s)");
let _ = tx.send(SweepMessage::Finished { dated });
});
});
rx
}
/// How many images still lack a date or a thumbnail.
pub(super) fn count_outstanding(catalog: &Catalog) -> Result<usize, dr_catalog::CatalogError> {
let n: i64 = catalog.connection().query_row(
&format!(
"SELECT count(*) FROM images
WHERE metadata_state < 2 AND {VISIBLE_UNALIASED}"
),
[],
|r| r.get(0),
)?;
Ok(n as usize)
}
/// The next images needing work.
///
/// Ordered by id so the sweep advances deterministically and a resumed run
/// picks up where it left off rather than revisiting.
pub(super) fn next_outstanding(
catalog: &Catalog,
limit: usize,
) -> Result<Vec<ThumbnailRequest>, dr_catalog::CatalogError> {
let mut stmt = catalog.connection().prepare(&format!(
"SELECT i.id, i.source_ref, r.file_id, i.file_size
FROM images i
LEFT JOIN remote r ON r.image_id = i.id
WHERE i.metadata_state < 2 AND {VISIBLE}
ORDER BY i.id
LIMIT ?1"
))?;
let rows = stmt
.query_map([limit as i64], |r| {
Ok(ThumbnailRequest {
// The sweep indexes dates, and reads headers only — the size
// never reaches a fetch, but it must name something.
thumb_size: dr_thumbs::ThumbSize::Grid,
// Row index is meaningless here — the sweep touches no grid
// cell, so nothing consumes it.
row: 0,
image_id: r.get(0)?,
path: r.get(1)?,
file_id: r.get::<_, Option<i64>>(2)?.map(|v| v as u64),
size: r.get::<_, Option<i64>>(3)?.unwrap_or(0) as u64,
needs_metadata: true,
full_resolution: false,
})
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
/// Commit a sweep chunk.
///
/// An image whose header yielded no date is still marked done, or the sweep
/// would revisit it forever. `write_metadata` records `metadata_state = 1` for
/// those, so this promotes them explicitly.
pub(super) fn flush_sweep(catalog: &Catalog, found: &mut Vec<MetadataFound>) {
if found.is_empty() {
return;
}
if let Err(e) = write_metadata(catalog, found) {
log::warn!("sweep: writing metadata: {e}");
found.clear();
return;
}
// Mark the dateless as examined. Without this they stay at state 1 and the
// sweep loops over them on every pass, never terminating.
let ids: Vec<i64> = found
.iter()
.filter(|m| m.captured_at.is_none())
.map(|m| m.image_id)
.collect();
for &id in &ids {
let _ = catalog
.connection()
.execute("UPDATE images SET metadata_state = 2 WHERE id = ?1", [id]);
}
// A header with no date is not a photograph with no date: WhatsApp and
// most re-exports strip EXIF and keep the date in the name.
if let Err(e) = dr_catalog::name_dates::fill(catalog.connection(), Some(&ids)) {
log::warn!("sweep: dating from names: {e}");
}
found.clear();
}
/// TRACES: FR-CULL-8 | FR-EXP-9
/// Open one original for a native render, orientation applied and nothing else.
///
/// The half of `repairs::detect` that has nothing to do with faces, exposed
/// because measuring what this pass is worth means rendering the same file two
/// ways and comparing the crops — see `examples/face_native.rs`. A tool that
/// had to reimplement the render would be measuring its own reimplementation.
pub fn render_native(gpu: &dr_gpu::GpuContext, bytes: &[u8]) -> Result<dr_export::Frame, String> {
open_native(gpu, dr_decode::default(), bytes)?.render_for_export(dr_types::ColourSpace::Srgb)
}
/// The session behind [`render_native`], crate-private because `DevelopSession`
/// is. The repair job needs the session itself rather than just its frame:
/// it renders the People screen's proxy from the same open session rather than
/// opening the file twice (`repairs::Fetched`).
pub(crate) fn open_native(
gpu: &dr_gpu::GpuContext,
decoder: &dyn dr_decode::Decoder,
bytes: &[u8],
) -> Result<crate::develop::DevelopSession, String> {
// Whatever the header says, or an empty one for a file that has none: the
// session takes its orientation from it, and remembers the rest for
// anything that later exports from this session (FR-EXP-8).
let meta = decoder.metadata(bytes).unwrap_or_default();
crate::open_session(gpu, decoder, bytes, &meta)
}
/// The class the whole-library pass fills.
///
/// Grid only, deliberately. The large class is four times the transfer for a
/// detail only a zoomed cell or the loupe asks for — on the reference library
/// that is ~200 MB of shards against ~860 MB, paid by *every* device that
/// syncs them (see [`dr_thumbs::ThumbSize`]). A photograph actually looked at
/// closely still gets its large thumbnail from the interactive path.
pub(super) const SWEEP_THUMB_SIZE: dr_thumbs::ThumbSize = dr_thumbs::ThumbSize::Grid;
/// Progress from the whole-library thumbnail pass.
#[derive(Debug)]
pub enum ThumbSweepMessage {
/// How many images still lack a thumbnail, counted once at the start.
Total(usize),
/// Another chunk finished. Carries cumulative counts.
Progress { done: usize, stored: usize },
Finished {
stored: usize,
failed: usize,
/// Stopped early because the server stopped answering. The pass is
/// resumable, so this is "come back later", not a failure.
offline: bool,
},
}
/// TRACES: FR-CAT-3 | FR-NC-3 | FR-NC-7
/// Thumbnail **every** image in the library, not just the ones browsed.
///
/// # Why this exists next to the grid's own fetching
///
/// The interactive path fills cells as they are scrolled past, which is the
/// right posture for a remote library (FR-NC-3) and the wrong one for handing
/// the result to a second device: a tablet that syncs the shards inherits only
/// the fraction of the library its sibling happened to look at. This is the
/// deliberate, user-launched version of the same work — an hour of range
/// fetches paid once, on the machine that can afford it, so every other client
/// gets a full grid for the cost of a few hundred MB (`derived_sync`).
///
/// # Shape, and why it borrows the metadata sweep's
///
/// Chunked and lane-parallel exactly as [`spawn_sweep`] is, for the same
/// reason: each image is ~0.6 s of round-trip latency and almost no
/// bandwidth, so the sequential version spends its life waiting. What differs
/// is the store — it is `&mut` and cannot be shared across lanes, which is why
/// the metadata sweep skips thumbnails entirely. Here the lanes fetch, decode
/// and *encode*, and only the ~20 KB result crosses back to this thread, which
/// owns the store and writes the chunk in one go. So the parallelism is real
/// and the single-writer rule is never bent.
///
/// Dates arrive free: the header a preview needs is the header EXIF lives in,
/// so an image this pass reaches is dated on the same fetch rather than
/// costing a second one.
///
/// Resumable by construction — the work list is what the store does not have,
/// so a kill costs the chunk in flight and nothing more. An image with no
/// locatable preview is retried on a later run; it is one header fetch, and
/// the alternative is a second piece of state that has to be invalidated when
/// a file is replaced.
pub fn spawn_thumbnail_sweep(
conn: Connection,
catalog_path: PathBuf,
store_dir: PathBuf,
) -> Receiver<ThumbSweepMessage> {
let (tx, rx) = std::sync::mpsc::channel();
// TRACES: FR-RAW-2
// The one place this job names a decoder; everything below takes it.
let decoder = dr_decode::default();
executors::spawn(Executor::Decode, "thumb-sweep", move || {
let finish_empty = |tx: &Sender<ThumbSweepMessage>| {
let _ = tx.send(ThumbSweepMessage::Finished {
stored: 0,
failed: 0,
offline: false,
});
};
let catalog = match Catalog::open(&catalog_path) {
Ok(c) => c,
Err(e) => {
log::warn!(
"thumbnail sweep: cannot open catalog at {}: {e}",
catalog_path.display()
);
finish_empty(&tx);
return;
}
};
// Unlike the grid's fetch, which carries on without a store and simply
// shows what it downloaded, a store that will not open ends this: the
// pass exists to fill it, and running an hour of transfers with
// nowhere to put them would be worse than not starting.
let mut store = match ThumbStore::open(&store_dir) {
Ok(s) => s,
Err(e) => {
log::warn!(
"thumbnail sweep: cannot open the thumbnail store at {}: {e}",
store_dir.display()
);
finish_empty(&tx);
return;
}
};
let wanted = match thumbnails_outstanding(&catalog, &store) {
Ok(w) => w,
Err(e) => {
log::warn!("thumbnail sweep: {e}");
finish_empty(&tx);
return;
}
};
let total = wanted.len();
if total == 0 {
log::info!("thumbnail sweep: every image already has a thumbnail");
finish_empty(&tx);
return;
}
log::info!("thumbnail sweep: {total} image(s) need a thumbnail");
if tx.send(ThumbSweepMessage::Total(total)).is_err() {
return;
}
let rt = match crate::net_runtime::build() {
Ok(rt) => rt,
Err(e) => {
log::warn!("thumbnail sweep: no runtime: {e}");
finish_empty(&tx);
return;
}
};
rt.block_on(async {
let backend = match crate::remote::connect(&conn) {
Ok(b) => b,
Err(e) => {
log::warn!("thumbnail sweep: {e}");
finish_empty(&tx);
return;
}
};
let (mut done, mut stored, mut failed) = (0usize, 0usize, 0usize);
let mut offline = false;
let mut found = Vec::new();
// TRACES: FR-NC-6c
// On a placeholder library the bytes may not be here at all, and
// this is a pass the user asked for — so it may fetch them, which
// browsing may not (ARCH §9.0a). Every file is *borrowed*: what
// this pass downloads it gives back, and what the user already had
// it leaves alone. Against a server or a plain folder every borrow
// is a no-op, so there is one code path rather than two.
let pool = dr_sync_folder::BorrowPool::new();
for chunk in wanted.chunks(SWEEP_CHUNK) {
// Each lane owns a disjoint slice and its own output, so
// nothing is shared and no lock is needed. The store is not
// touched here — see the note on the function.
let lanes: Vec<Vec<&ThumbnailRequest>> = (0..SWEEP_LANES)
.map(|lane| chunk.iter().skip(lane).step_by(SWEEP_LANES).collect())
.collect();
let results = futures_join_all(lanes.into_iter().map(|lane| {
let backend: &dyn RemoteBackend = &*backend;
let pool = &pool;
async move {
let mut made: Vec<(u64, dr_thumbs::Thumbnail)> = Vec::new();
let mut found = Vec::new();
let mut attempted = 0usize;
let mut failed = 0usize;
let mut offline = false;
for req in lane {
// Enforced by the query, which joins `remote`: an
// image with no file id has nothing to key the
// store on and is not a candidate.
let Some(file_id) = req.file_id else { continue };
attempted += 1;
// Held for this image only. A failure to fetch is
// this image's verdict, not the batch's: a client
// that cannot reach the server reports it as
// offline through the usual path below.
let _held =
match pool.borrow(backend, &RemotePath::new(&req.path)).await {
Ok(h) => h,
Err(e) if e.indicates_offline() => {
log::info!("thumbnail sweep: {e}");
attempted -= 1;
offline = true;
break;
}
Err(e) => {
log::debug!("thumbnail sweep: {}: {e}", req.path);
failed += 1;
continue;
}
};
match fetch_preview(backend, decoder, req, &mut found).await {
PreviewOutcome::Ready(preview) => {
match encode_preview(file_id, &preview) {
Some(thumb) => made.push((file_id, thumb)),
None => failed += 1,
}
}
PreviewOutcome::Unavailable(reason) => {
log::debug!("thumbnail sweep: {}: {reason}", req.path);
failed += 1;
}
// Nothing after this would reach the server
// either, so the lane stops rather than
// spending a timeout per remaining image.
PreviewOutcome::Offline(reason) => {
log::info!("thumbnail sweep: server unreachable: {reason}");
attempted -= 1;
offline = true;
break;
}
}
}
(made, found, attempted, failed, offline)
}
}))
.await;
for (made, lane_found, attempted, lane_failed, lane_offline) in results {
done += attempted;
failed += lane_failed;
offline |= lane_offline;
found.extend(lane_found);
for (file_id, thumb) in made {
if store_thumbnail(&mut store, file_id, SWEEP_THUMB_SIZE, &thumb) {
stored += 1;
} else {
failed += 1;
}
}
}
// Committed per chunk rather than at the end, so a kill keeps
// every date read so far — the same bargain the metadata sweep
// makes, and for the same reason.
flush_sweep(&catalog, &mut found);
if tx
.send(ThumbSweepMessage::Progress { done, stored })
.is_err()
{
// Cancelled. Hand back what was borrowed before leaving,
// or a stopped pass costs the disk of everything it had
// reached and delivers nothing for it.
pool.release_all(&*backend).await;
return;
}
if offline {
break;
}
}
flush_sweep(&catalog, &mut found);
// Give back everything this pass fetched, before reporting done —
// a user watching the disk should see it return, and a pass that
// reported success while still holding the library would be
// lying about what it cost.
let returned = pool.release_all(&*backend).await;
if returned.released > 0 {
log::info!(
"thumbnail sweep: released {} borrowed file(s)",
returned.released
);
}
log::info!("thumbnail sweep: {stored} stored, {failed} without a usable preview");
let _ = tx.send(ThumbSweepMessage::Finished {
stored,
failed,
offline,
});
});
});
rx
}
/// Every visible image on the server that the store has no grid thumbnail for.
///
/// Joined against `remote` rather than left-joined: the store is keyed on
/// Nextcloud's `oc:fileid` (FR-NC-5), so an image the scan recorded without
/// one cannot be stored and is not work this pass can do.
///
/// The whole list is built up front rather than re-queried per chunk, unlike
/// the metadata sweep: "does the store have this" is answered by the store's
/// index, which this thread is also the one writing, so a stale list is not a
/// risk the way a concurrently-dating grid made it one there.
pub(super) fn thumbnails_outstanding(
catalog: &Catalog,
store: &ThumbStore,
) -> Result<Vec<ThumbnailRequest>, dr_catalog::CatalogError> {
let mut stmt = catalog.connection().prepare(&format!(
"SELECT i.id, i.source_ref, r.file_id, i.file_size, i.metadata_state
FROM images i
JOIN remote r ON r.image_id = i.id
WHERE r.file_id IS NOT NULL AND {VISIBLE}
ORDER BY i.id"
))?;
let rows = stmt
.query_map([], |r| {
let file_id = r.get::<_, Option<i64>>(2)?.map(|v| v as u64);
Ok(ThumbnailRequest {
thumb_size: SWEEP_THUMB_SIZE,
// No grid cell is waiting on this, so nothing consumes the row.
row: 0,
image_id: r.get(0)?,
path: r.get(1)?,
file_id,
size: r.get::<_, Option<i64>>(3)?.unwrap_or(0) as u64,
// The header this fetch reads is the one EXIF lives in, so an
// undated image is dated on the way past for nothing.
needs_metadata: r.get::<_, i64>(4)? < 2,
full_resolution: false,
})
})?
.filter_map(Result::ok)
.filter(|req| {
req.file_id
.is_some_and(|id| !store.contains(id, SWEEP_THUMB_SIZE))
})
.collect();
Ok(rows)
}
#[cfg(test)]
mod tests {
use super::*;
/// A backend whose every read fails the same way, counting the attempts.
struct Refusing {
error: fn() -> dr_sync::RemoteError,
gets: std::sync::atomic::AtomicUsize,
caps: dr_sync::Capabilities,
}
#[async_trait::async_trait]
impl RemoteBackend for Refusing {
fn capabilities(&self) -> &dr_sync::Capabilities {
&self.caps
}
fn name(&self) -> &str {
"refusing"
}
async fn list(
&self,
_dir: &RemotePath,
_since: Option<&dr_sync::Validator>,
) -> Result<Vec<dr_sync::RemoteEntry>, dr_sync::RemoteError> {
Err((self.error)())
}
async fn dir_validator(
&self,
_dir: &RemotePath,
) -> Result<dr_sync::Validator, dr_sync::RemoteError> {
Err((self.error)())
}
async fn delta(
&self,
_c: &dr_sync::Cursor,
) -> Result<(Vec<dr_sync::RemoteChange>, dr_sync::Cursor), dr_sync::RemoteError> {
Err((self.error)())
}
async fn get(
&self,
_id: &RemoteId,
_r: Option<std::ops::Range<u64>>,
) -> Result<Vec<u8>, dr_sync::RemoteError> {
self.gets.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Err((self.error)())
}
async fn put(
&self,
_p: &RemotePath,
_b: Vec<u8>,
_c: Option<dr_sync::Precondition>,
) -> Result<dr_sync::Validator, dr_sync::RemoteError> {
Err((self.error)())
}
async fn delete(
&self,
_id: &RemoteId,
_c: Option<dr_sync::Precondition>,
) -> Result<(), dr_sync::RemoteError> {
Err((self.error)())
}
async fn move_to(
&self,
_f: &RemoteId,
_t: &RemotePath,
) -> Result<(), dr_sync::RemoteError> {
Err((self.error)())
}
async fn create_dir(&self, _p: &RemotePath) -> Result<(), dr_sync::RemoteError> {
Err((self.error)())
}
}
fn read_date_against(error: fn() -> dr_sync::RemoteError) -> (DateRead, usize) {
let backend = Refusing {
error,
gets: Default::default(),
caps: dr_sync::Capabilities::minimal(),
};
let req = ThumbnailRequest {
row: 0,
path: "a.CR2".into(),
file_id: Some(1),
size: 0,
image_id: 1,
thumb_size: dr_thumbs::ThumbSize::Grid,
needs_metadata: true,
full_resolution: false,
};
let rt = crate::net_runtime::build().unwrap();
let outcome = rt.block_on(read_metadata_only(
&backend,
dr_decode::default(),
&req,
&mut Vec::new(),
));
(
outcome,
backend.gets.load(std::sync::atomic::Ordering::SeqCst),
)
}
#[test]
fn a_date_read_against_a_dead_server_asks_once_and_says_so() {
// Each attempt waited out a 15 s connect timeout, three per image, for
// every cached-but-undated cell in the window — minutes of "reading
// dates" with no banner, because the caller could not tell a dead
// server from a missing file.
let (outcome, gets) =
read_date_against(|| dr_sync::RemoteError::Network("connection refused".into()));
assert!(matches!(outcome, DateRead::Offline(_)), "{outcome:?}");
assert_eq!(gets, 1, "retrying a dead server buys another timeout");
}
#[test]
fn a_date_read_that_hits_a_lock_still_retries() {
// The reason the retry exists: Nextcloud answers a read with 423 under
// our own concurrency, and the same range succeeds moments later.
let (outcome, gets) = read_date_against(|| dr_sync::RemoteError::Server {
status: 423,
detail: "locked".into(),
});
assert_eq!(outcome, DateRead::Failed);
assert_eq!(gets, 3);
}
#[test]
fn a_header_gives_the_size_the_photograph_is_seen_at() {
let mut md = dr_decode::Metadata {
width: Some(9000),
height: Some(3000),
..Default::default()
};
assert_eq!(MetadataFound::upright_size(&md), Some((9000, 3000)));
md.orientation = Some(dr_types::Orientation::from_exif(6));
assert_eq!(MetadataFound::upright_size(&md), Some((3000, 9000)));
md.width = None;
assert_eq!(MetadataFound::upright_size(&md), None);
}
#[test]
fn join_all_preserves_order_regardless_of_completion() {
// The ordering guarantee is what lets a caller pair results back to
// their inputs; without it a lane's dates could be attributed to the
// wrong images.
let rt = crate::net_runtime::build().unwrap();
let out = rt.block_on(async {
futures_join_all(vec![
Box::pin(async { 1 }) as std::pin::Pin<Box<dyn std::future::Future<Output = i32>>>,
Box::pin(async {
tokio::task::yield_now().await;
tokio::task::yield_now().await;
2
}),
Box::pin(async {
tokio::task::yield_now().await;
3
}),
])
.await
});
assert_eq!(out, vec![1, 2, 3]);
}
#[test]
fn join_all_of_nothing_completes() {
let rt = crate::net_runtime::build().unwrap();
let out: Vec<i32> =
rt.block_on(async { futures_join_all(Vec::<std::future::Ready<i32>>::new()).await });
assert!(out.is_empty());
}
#[test]
fn sweep_lanes_divide_a_chunk_without_loss() {
// Every image in a chunk must land in exactly one lane: a striding
// split that dropped or duplicated one would silently under- or
// double-index the library.
let chunk: Vec<usize> = (0..SWEEP_CHUNK).collect();
let lanes: Vec<Vec<usize>> = (0..SWEEP_LANES)
.map(|l| chunk.iter().skip(l).step_by(SWEEP_LANES).copied().collect())
.collect();
let mut seen: Vec<usize> = lanes.iter().flatten().copied().collect();
seen.sort_unstable();
assert_eq!(seen, chunk);
// Evenly divided, so no lane sits idle while another finishes.
assert!(lanes.iter().all(|l| l.len() == SWEEP_CHUNK / SWEEP_LANES));
}
#[test]
fn a_short_chunk_still_covers_every_image() {
// The last chunk of a library is rarely a full multiple of the lanes.
let chunk: Vec<usize> = (0..5).collect();
let lanes: Vec<Vec<usize>> = (0..SWEEP_LANES)
.map(|l| chunk.iter().skip(l).step_by(SWEEP_LANES).copied().collect())
.collect();
let mut seen: Vec<usize> = lanes.iter().flatten().copied().collect();
seen.sort_unstable();
assert_eq!(seen, chunk);
}
}