//! 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, pub captured_offset: Option, pub camera: Option, pub lens: Option, pub iso: Option, /// 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, tx: &Sender, ) { 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, ) -> 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 { 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(futures: impl IntoIterator) -> Vec 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>>> = futures.into_iter().map(|f| Some(Box::pin(f))).collect(); let mut done: Vec> = (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 { 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> = (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 = 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 { 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, 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>(2)?.map(|v| v as u64), size: r.get::<_, Option>(3)?.unwrap_or(0) as u64, needs_metadata: true, full_resolution: false, }) })? .collect::, _>>()?; 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) { 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 = 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 { 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 { // 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 { 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| { 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> = (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, 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>(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>(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 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::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 = rt.block_on(async { futures_join_all(Vec::>::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 = (0..SWEEP_CHUNK).collect(); let lanes: Vec> = (0..SWEEP_LANES) .map(|l| chunk.iter().skip(l).step_by(SWEEP_LANES).copied().collect()) .collect(); let mut seen: Vec = 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 = (0..5).collect(); let lanes: Vec> = (0..SWEEP_LANES) .map(|l| chunk.iter().skip(l).step_by(SWEEP_LANES).copied().collect()) .collect(); let mut seen: Vec = lanes.iter().flatten().copied().collect(); seen.sort_unstable(); assert_eq!(seen, chunk); } }