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.
940 lines
38 KiB
Rust
940 lines
38 KiB
Rust
//! 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>,
|
||
}
|
||
|
||
/// 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.
|
||
///
|
||
/// Returns 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>,
|
||
) -> bool {
|
||
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 true;
|
||
}
|
||
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 false;
|
||
}
|
||
}
|
||
}
|
||
false
|
||
}
|
||
|
||
/// 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)
|
||
WHERE id = ?1",
|
||
rusqlite::params![
|
||
m.image_id,
|
||
m.captured_at,
|
||
m.captured_offset,
|
||
m.camera,
|
||
m.lens,
|
||
m.iso,
|
||
state,
|
||
],
|
||
)?;
|
||
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 {
|
||
if read_metadata_only(backend, decoder, req, &mut found).await {
|
||
reached.push(req.image_id);
|
||
}
|
||
}
|
||
(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,
|
||
}),
|
||
);
|
||
|
||
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]);
|
||
}
|
||
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::*;
|
||
|
||
#[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);
|
||
}
|
||
}
|