From 846a2491569090120c67a276629ad9086f47e711 Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Sat, 29 Aug 2026 23:31:47 +0200 Subject: [PATCH] Drain the queue that nothing has ever drained MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `jobs` has been a complete durable work queue since the catalog was written, and nothing has ever taken a job out of it. `claim_next`, `complete`, `fail` and `recover_orphaned` had no callers outside their own tests; `enqueue` had three. So the table grew one row per photograph and kept it forever, and FR-PLAT-AND-3's resumability was a property of code that never ran. `runner` is the missing half. It owns no thread, no clock and no policy, and that is the whole design: on Android the process does not decide when background work may run. WorkManager does, subject to Doze, battery saver and FR-NC-6's network constraints, and it revokes permission mid-job by calling onStopped(). So the runner exposes `run_one` — claim, run, record — and `drain`, which repeats it against a budget, a deadline and a cancellation flag the host owns. A `Worker.doWork()` with ten minutes calls drain with a deadline; a desktop idle pass calls it with none. That is the seam the Android service plugs into, and it needs no Android to test. Handlers are supplied from above, because the catalog knows what needs doing and nothing about how: a thumbnail needs a decoder and a fetch needs a network stack, neither of which belongs under core/dr-catalog. A runner claims only kinds some handler declares, so a queue holding work this device cannot do is left alone rather than failed five times. Four outcomes, and only two of them are the job's fault. Done deletes the row; Retry backs off; Abandon gives up now, for a failure no retry can fix; Interrupted releases the claim with its attempt refunded and ends the drain, because the host stopped rather than the job — five backgroundings in a row must not mark good work as failed. Process death is the fifth and cannot report itself, which is what `recover` is for. Recovery is called from `show_catalog_now`, which is the one place a catalog is opened for a session and already returns early if one is open. It has to be exactly once and before any worker starts: there is no owner column, so a second pass while a worker held a claim would take it away. The attempt a dead claim consumed is deliberately kept — a job that takes the process down with it is indistinguishable from one that fails, and the attempt counter is the only evidence that survives a death. The tests cover claiming under contention twice over: sequentially across two connections, and with four threads on four connections against one catalog on disk, asserting every job ran exactly once. Plus completion, backoff, giving up, abandoning, interruption, budget, deadline, cancellation, and a job orphaned by a simulated crash being reclaimed and run once rather than lost or repeated. Not wired to a handler yet, and deliberately not: the only enqueue site the app actually reaches is the remote scan's, whose thumbnails are already served by the async grid worker, and `walk`'s two sites are reachable only from the scan_local example. Inventing a handler to make the plumbing look used is how a requirement comes to read as covered by code that does not implement it. Co-Authored-By: Claude Opus 5 (1M context) --- core/dr-catalog/src/lib.rs | 6 + core/dr-catalog/src/runner.rs | 930 ++++++++++++++++++++++++++++++++++ ui/dr-ui/src/library_ui.rs | 23 + 3 files changed, 959 insertions(+) create mode 100644 core/dr-catalog/src/runner.rs diff --git a/core/dr-catalog/src/lib.rs b/core/dr-catalog/src/lib.rs index a638cc6..0ca953f 100644 --- a/core/dr-catalog/src/lib.rs +++ b/core/dr-catalog/src/lib.rs @@ -18,6 +18,7 @@ //! - [`faces`] — detected faces, the people they belong to, and who said so //! - [`bursts`] — frames that are one moment, grouped so they judge as one //! - [`jobs`] — the durable background work queue +//! - [`runner`] — the thing that drains it, driven by whoever owns the thread //! - [`trash`] — soft delete to a folder, then permanent delete //! - [`merge`] / [`sync`] — cross-device merging of collections and keywords //! @@ -46,6 +47,7 @@ pub mod keywords; pub mod merge; pub mod query; pub mod rating; +pub mod runner; pub mod scan; pub mod schema; pub mod sync; @@ -63,6 +65,10 @@ pub use keywords::{Coverage, Keyword, KeywordId, SelectionKeyword}; pub use merge::MergeReport; pub use query::{Query, Sort}; pub use rating::{Judgement, MAX_RATING}; +// Not `runner::Budget`: `cache::Budget` already owns that name here and +// means something else entirely (bytes on disk, not jobs in a slot). +// Callers spell the work budget `runner::Budget`, where it is unambiguous. +pub use runner::{DrainReport, JobHandler, Outcome, Runner}; pub use scan::{DirAction, DirState, EntryAction, ScanOutcome}; pub use trash::{TrashedImage, TRASH_DIR}; pub use walk::{ensure_root, mark_root_offline, scan_root, RootKind, ScanProgress, ScanReport}; diff --git a/core/dr-catalog/src/runner.rs b/core/dr-catalog/src/runner.rs new file mode 100644 index 0000000..b47f850 --- /dev/null +++ b/core/dr-catalog/src/runner.rs @@ -0,0 +1,930 @@ +//! TRACES: FR-PLAT-AND-4 | FR-PLAT-AND-3 +//! The thing that drains the queue. +//! +//! [`crate::jobs`] has been a complete, durable, coalescing work queue since +//! the catalog was written, and nothing has ever taken a job out of it. Every +//! producer — the local walk, the remote scan — called `enqueue` and no one +//! called `claim_next`, so the table grew one row per photograph and stayed +//! that size forever. This module is the missing half. +//! +//! # Why the runner is driven rather than self-owning +//! +//! The obvious shape is a thread that loops until the queue is empty, and it +//! is the wrong one. On Android the process does not decide when background +//! work may run: `WorkManager` does, subject to Doze, battery saver and the +//! metered-network constraints in FR-NC-6, and it revokes permission mid-job +//! by calling `onStopped()` (FR-PLAT-AND-4). A foreground service for a +//! user-initiated export gets a longer leash but still not an unbounded one. +//! +//! So the runner owns no thread, no clock and no policy. It exposes +//! [`Runner::run_one`] — claim one job, run it, record what happened — and +//! [`Runner::drain`], which repeats that against a [`Budget`] and a +//! cancellation flag the host owns. A `Worker.doWork()` that must return +//! within ten minutes calls `drain` with a deadline; a desktop idle loop calls +//! it with none. Neither has to reach inside. +//! +//! Everything the host supplies is passed in for the same reason `jobs` takes +//! `now` rather than reading the clock: a scheduler is exactly the thing that +//! has to be testable without waiting. +//! +//! # Why interruption is not failure +//! +//! Four things can happen to a claimed job, and only two of them are the job's +//! fault: +//! +//! - [`Outcome::Done`] — the row is deleted. +//! - [`Outcome::Retry`] — the work failed and might succeed later. Backoff, +//! and eventually [`crate::jobs::MAX_ATTEMPTS`] gives up on it. +//! - [`Outcome::Abandon`] — the work cannot succeed, ever. Failing five times +//! over five minutes to learn that is five minutes of a phone's battery. +//! - [`Outcome::Interrupted`] — the *host* stopped, not the job. The claim is +//! released and the attempt it consumed is given back, because a user who +//! pulled the app off the screen has not told us anything about the file. +//! +//! Process death is the fifth case and the one that cannot report itself: the +//! row simply stays `Running` with no owner. [`Runner::recover`] is what +//! reclaims it, and it is why an interrupted job is resumable rather than lost +//! (FR-PLAT-AND-3). It must run **before** any worker starts against a +//! catalog, or it will steal a job another runner is holding — there is no +//! owner column to tell them apart. + +use std::sync::atomic::{AtomicBool, Ordering}; + +use rusqlite::Connection; + +use crate::error::CatalogError; +use crate::jobs::{self, Job, JobKind}; + +/// What running a job turned out to be. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum Outcome { + /// The work is done. The row goes away. + Done, + /// It failed, and trying again later is reasonable. Backoff applies, and + /// [`crate::jobs::MAX_ATTEMPTS`] eventually stops it. + Retry(String), + /// It failed in a way no retry can fix — the subject is gone, the payload + /// is unreadable, the format is one this build does not know. Marked + /// failed at once rather than burning the whole retry ladder to reach the + /// same answer. + Abandon(String), + /// The host is stopping, and the job never really ran. + /// + /// Distinct from `Retry` because it costs no attempt: `onStopped()` five + /// times in a row would otherwise mark a perfectly good job as failed. + Interrupted, +} + +/// Something that can actually do the work a job describes. +/// +/// The catalog knows what needs doing and nothing about how — a thumbnail +/// needs a decoder, a fetch needs a network stack, and neither belongs under +/// `core/dr-catalog` (ARCH §4.1: calls go downward). So the queue lives here +/// and the handlers are supplied from above. +pub trait JobHandler { + /// The kinds this handler will accept. + /// + /// Load-bearing, not documentation: the runner claims **only** kinds some + /// handler declares. A queue holding `FetchOriginal` rows on a device with + /// no connector must leave them alone rather than claim them and fail + /// them, and a runner that claimed everything would do exactly that — five + /// times each, with backoff, on battery. + fn kinds(&self) -> &[JobKind]; + + /// Do the work. + /// + /// The connection is offered because most handlers write their result back + /// into the catalog; one that does not is free to ignore it. It is the + /// runner's own connection, so a handler must not hold a transaction open + /// across a network call — the runner needs it back to record the outcome. + fn run(&mut self, conn: &Connection, job: &Job) -> Outcome; +} + +/// How much work a host is willing to let one drain do. +/// +/// Both limits are checked *before* a job is claimed, never during one: a +/// handler is opaque and may be halfway through writing a sidecar. Overrunning +/// a deadline by one job is survivable; being killed mid-write is the thing +/// [`Runner::recover`] exists to clean up after. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub struct Budget { + /// Stop after this many jobs. `None` means "until the queue is empty". + pub max_jobs: Option, + /// Stop once the clock reaches this second. Same clock the drain is given. + pub deadline: Option, +} + +impl Budget { + /// Run until nothing is left. What a desktop idle pass wants. + pub const UNLIMITED: Self = Self { + max_jobs: None, + deadline: None, + }; + + /// At most `n` jobs. A slice small enough to stay responsive. + pub fn jobs(n: usize) -> Self { + Self { + max_jobs: Some(n), + ..Self::UNLIMITED + } + } + + /// Until the clock reaches `deadline`. What a `WorkManager` slot wants. + pub fn until(deadline: i64) -> Self { + Self { + deadline: Some(deadline), + ..Self::UNLIMITED + } + } +} + +/// Why a drain stopped. +/// +/// Worth distinguishing because the host's next move differs: `Drained` means +/// there is nothing to reschedule for, and the other three all mean "there is +/// more, ask again". +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub enum Stopped { + /// Nothing claimable is left. + #[default] + Drained, + /// The job count ran out. + Budget, + /// The clock ran out. + Deadline, + /// The host asked it to stop, or a handler reported itself interrupted. + Cancelled, +} + +impl Stopped { + /// Whether the queue may still hold claimable work. + pub fn more_to_do(self) -> bool { + self != Stopped::Drained + } +} + +/// What one drain did. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub struct DrainReport { + pub completed: usize, + pub retried: usize, + pub abandoned: usize, + pub interrupted: usize, + pub stopped: Stopped, +} + +impl DrainReport { + /// Jobs claimed, whatever became of them. This is what a budget counts. + pub fn ran(&self) -> usize { + self.completed + self.retried + self.abandoned + self.interrupted + } +} + +/// What a recovery pass found waiting from the last run. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub struct Recovered { + /// Jobs a dead process was holding. These are the resumed ones. + pub reclaimed: usize, + /// Jobs deleted because the photograph they name no longer exists. + pub reaped: usize, +} + +impl Recovered { + pub fn did_anything(&self) -> bool { + self.reclaimed > 0 || self.reaped > 0 + } +} + +/// Ready the queue for a fresh run, before any worker touches it. +/// +/// Two distinct cleanups, and both are startup-only: +/// +/// - **Reclaim.** A `Running` row has no owner; the process that claimed it is +/// gone. On Android that is a routine morning, not a crash (FR-PLAT-AND-3). +/// The attempt it consumed is *kept*, deliberately: a job that takes the +/// process down with it three times running should not be retried forever, +/// and the attempt counter is the only evidence of that we have. +/// - **Reap.** Jobs naming an image the catalog no longer has. A library that +/// has been culled leaves thumbnail jobs for photographs that were deleted +/// months ago, and every one of them would be claimed, run and failed. +/// +/// Reclaim runs first so its count is the honest number of interrupted jobs, +/// before reaping removes whichever of them pointed at nothing. +/// +/// **Call this exactly once per catalog, at startup.** It cannot distinguish a +/// job a dead process was holding from one a live runner is holding right now, +/// because there is no owner column — the queue is durable, not distributed. +pub fn recover(conn: &Connection) -> Result { + Ok(Recovered { + reclaimed: jobs::recover_orphaned(conn)?, + reaped: jobs::reap_orphan_subjects(conn)?, + }) +} + +/// Claims work, runs it, and records what happened. +/// +/// Borrows its connection rather than owning one so a host can drive it from +/// the same handle it already has open. Nothing here spawns a thread; several +/// runners on several threads, each with its own connection to the same +/// catalog, are safe because the claim is a single atomic statement (see +/// [`crate::jobs::claim_next`]). +pub struct Runner<'a> { + conn: &'a Connection, + handlers: Vec>, + /// The union of every handler's kinds, cached because it is passed to + /// every claim. This is what stops the runner claiming work it cannot do. + claimable: Vec, +} + +impl<'a> Runner<'a> { + /// A runner with no handlers. It can recover, and it can claim nothing. + pub fn new(conn: &'a Connection) -> Self { + Self { + conn, + handlers: Vec::new(), + claimable: Vec::new(), + } + } + + /// Add a handler. + pub fn with(self, handler: impl JobHandler + 'a) -> Self { + self.with_boxed(Box::new(handler)) + } + + /// Add a handler chosen at runtime — a network one only where there is a + /// connector, a decoding one only where there is a decoder. + pub fn with_boxed(mut self, handler: Box) -> Self { + for kind in handler.kinds() { + if !self.claimable.contains(kind) { + self.claimable.push(*kind); + } + } + self.handlers.push(handler); + self + } + + /// The kinds this runner will claim. Useful to a host deciding whether + /// starting it is worth waking the radio for. + pub fn claimable(&self) -> &[JobKind] { + &self.claimable + } + + /// See [`recover`]. Offered here too so a host has one thing to hold. + pub fn recover(&self) -> Result { + recover(self.conn) + } + + /// Claim one job, run it, and record the outcome. + /// + /// `Ok(None)` means nothing this runner can do is claimable *now* — the + /// queue may still hold work of other kinds, or work still backing off. + pub fn run_one(&mut self, now: i64) -> Result, CatalogError> { + // Copied out before `self.handlers` is borrowed mutably below. Both + // are fields of `self`, but the copy is what lets the two borrows + // coexist without the connection being reborrowed through `self`. + let conn = self.conn; + + let Some(job) = jobs::claim_next_matching(conn, now, &self.claimable)? else { + return Ok(None); + }; + + let outcome = match self + .handlers + .iter_mut() + .find(|h| h.kinds().contains(&job.kind)) + { + Some(handler) => handler.run(conn, &job), + // Unreachable by construction: `claimable` is exactly the union of + // the handlers' kinds. Parked rather than released, because + // releasing it would put it straight back where the next turn of + // the drain loop would claim it again, forever. + None => Outcome::Abandon(format!("no handler for {:?}", job.kind)), + }; + + match &outcome { + Outcome::Done => jobs::complete(conn, job.id)?, + Outcome::Retry(why) => jobs::fail(conn, &job, now, why)?, + Outcome::Abandon(why) => jobs::abandon(conn, job.id, why)?, + Outcome::Interrupted => jobs::release(conn, &job)?, + } + + Ok(Some(Ran { job, outcome })) + } + + /// Run jobs until the budget, the flag or the queue says stop. + /// + /// `clock` is called once per iteration rather than sampled once, because + /// the two things it feeds both move during a long drain: the deadline + /// check, and the `now` a failing job's backoff is measured from. + /// + /// Cancellation is checked between jobs only. A handler that wants to bail + /// out of work already started says so with [`Outcome::Interrupted`], + /// which also ends the drain — otherwise a handler that always interrupts + /// would release its job and be handed it straight back. + pub fn drain( + &mut self, + clock: &dyn Fn() -> i64, + budget: Budget, + cancel: &AtomicBool, + ) -> Result { + let mut report = DrainReport::default(); + + loop { + // Relaxed: the flag is a one-way latch set by another thread and + // the only thing ordered against it is our own next claim. Missing + // one turn of the loop costs a job, not correctness. + if cancel.load(Ordering::Relaxed) { + report.stopped = Stopped::Cancelled; + break; + } + if budget.max_jobs.is_some_and(|max| report.ran() >= max) { + report.stopped = Stopped::Budget; + break; + } + + let now = clock(); + if budget.deadline.is_some_and(|end| now >= end) { + report.stopped = Stopped::Deadline; + break; + } + + let Some(ran) = self.run_one(now)? else { + report.stopped = Stopped::Drained; + break; + }; + + match ran.outcome { + Outcome::Done => report.completed += 1, + Outcome::Retry(_) => report.retried += 1, + Outcome::Abandon(_) => report.abandoned += 1, + Outcome::Interrupted => { + report.interrupted += 1; + report.stopped = Stopped::Cancelled; + break; + } + } + } + + Ok(report) + } + + /// Drain with no budget and no cancellation, at a fixed instant. + /// + /// Terminates because a job that fails is pushed past `now` by its backoff + /// and stops being claimable at this instant. + pub fn drain_all(&mut self, now: i64) -> Result { + static NEVER: AtomicBool = AtomicBool::new(false); + self.drain(&|| now, Budget::UNLIMITED, &NEVER) + } +} + +/// One job and what became of it. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Ran { + pub job: Job, + pub outcome: Outcome, +} + +#[cfg(test)] +mod tests { + use std::path::{Path, PathBuf}; + use std::sync::{Arc, Mutex}; + + use super::*; + use crate::jobs::{enqueue, JobState, Priority, MAX_ATTEMPTS}; + use crate::schema; + + /// A handler built from a closure, so each test states its own behaviour. + struct Fake { + kinds: Vec, + act: F, + } + + impl Outcome> JobHandler for Fake { + fn kinds(&self) -> &[JobKind] { + &self.kinds + } + fn run(&mut self, _conn: &Connection, job: &Job) -> Outcome { + (self.act)(job) + } + } + + fn handler Outcome>(kinds: &[JobKind], act: F) -> Fake { + Fake { + kinds: kinds.to_vec(), + act, + } + } + + /// A handler that records which subjects it saw and always succeeds. + /// + /// Takes its kinds by value and borrows nothing, so the returned handler is + /// `Send + 'static` and can be moved into a worker thread — which the + /// contention test needs. + fn recording(kinds: Vec, seen: Arc>>) -> impl JobHandler + Send { + Fake { + kinds, + act: move |job: &Job| { + seen.lock().unwrap().push(job.subject_id.unwrap_or(-1)); + Outcome::Done + }, + } + } + + fn db() -> Connection { + let c = Connection::open_in_memory().unwrap(); + schema::configure(&c).unwrap(); + schema::migrate(&c).unwrap(); + c + } + + /// An image row, so a job has a subject that exists. + fn image(c: &Connection, id: i64) { + c.execute( + "INSERT INTO roots(id, kind, label) VALUES (1, 'local', '/lib') + ON CONFLICT DO NOTHING", + [], + ) + .unwrap(); + c.execute( + "INSERT INTO images(id, root_id, source_ref, added_at) VALUES (?1, 1, ?2, 0)", + rusqlite::params![id, format!("/lib/{id}.CR3")], + ) + .unwrap(); + } + + fn queued(c: &Connection, kind: JobKind, subject: i64) { + image(c, subject); + enqueue(c, kind, Some(subject), Priority::Background, None).unwrap(); + } + + fn rows(c: &Connection) -> i64 { + c.query_row("SELECT count(*) FROM jobs", [], |r| r.get(0)) + .unwrap() + } + + /// A catalog on disk, so more than one connection can open it. + fn temp_catalog(name: &str) -> PathBuf { + let dir = std::env::temp_dir().join(format!( + "dr-runner-{name}-{}-{:?}", + std::process::id(), + std::thread::current().id() + )); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + dir.join("catalog.db") + } + + fn open(path: &Path) -> Connection { + let c = Connection::open(path).unwrap(); + schema::configure(&c).unwrap(); + schema::migrate(&c).unwrap(); + // Several connections write to this file at once in the contention + // tests. Without a busy handler the loser of a race gets an error + // instead of a turn. + c.busy_timeout(std::time::Duration::from_secs(10)).unwrap(); + c + } + + #[test] + fn a_completed_job_leaves_the_queue() { + // The whole finding in one assertion: before this module, the row + // stayed forever because nothing ever claimed it. + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + + let seen = Arc::new(Mutex::new(Vec::new())); + let report = Runner::new(&c) + .with(recording(vec![JobKind::Thumbnail], seen.clone())) + .drain_all(0) + .unwrap(); + + assert_eq!(report.completed, 1); + assert_eq!(report.stopped, Stopped::Drained); + assert_eq!(*seen.lock().unwrap(), vec![1]); + assert_eq!(rows(&c), 0, "a completed job leaves no row behind"); + } + + #[test] + fn a_failed_job_backs_off_and_is_claimed_again_later() { + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + + // A `Cell` rather than a captured `bool`, so the closure's mutability + // is its own business and the test reads the same either way. + let failed_once = std::cell::Cell::new(false); + let mut runner = Runner::new(&c).with(handler(&[JobKind::Thumbnail], |_| { + if failed_once.replace(true) { + Outcome::Done + } else { + Outcome::Retry("decoder said no".into()) + } + })); + + let first = runner.drain_all(100).unwrap(); + assert_eq!(first.retried, 1); + assert_eq!(rows(&c), 1, "a retryable failure keeps its row"); + + // Still inside the backoff window: nothing claimable, so the drain + // reports itself drained rather than spinning on the same job. + assert_eq!(runner.drain_all(100).unwrap().ran(), 0); + + let later = runner.drain_all(100 + jobs::backoff_seconds(1)).unwrap(); + assert_eq!(later.completed, 1); + assert_eq!(rows(&c), 0); + } + + #[test] + fn a_job_that_keeps_failing_is_given_up_on() { + // FR-RAW-4: one corrupt file must not stall the queue behind endless + // retries. Driven through the runner rather than by hand, because the + // runner is what a corrupt file will actually meet. + let c = db(); + queued(&c, JobKind::ExtractMetadata, 1); + + let mut runner = Runner::new(&c).with(handler(&[JobKind::ExtractMetadata], |_| { + Outcome::Retry("corrupt file".into()) + })); + + let mut now = 0; + for _ in 0..MAX_ATTEMPTS { + assert_eq!(runner.drain_all(now).unwrap().retried, 1); + now += jobs::backoff_seconds(MAX_ATTEMPTS); + } + + assert_eq!(runner.drain_all(now + 100_000).unwrap().ran(), 0); + let state: i64 = c + .query_row("SELECT state FROM jobs", [], |r| r.get(0)) + .unwrap(); + assert_eq!(state, JobState::Failed as i64); + } + + #[test] + fn an_abandoned_job_is_not_retried_at_all() { + // The difference that matters on battery: five failures spread over + // five minutes to learn what the first one already said. + let c = db(); + queued(&c, JobKind::FetchOriginal, 1); + + let mut runner = Runner::new(&c).with(handler(&[JobKind::FetchOriginal], |_| { + Outcome::Abandon("no connector on this device".into()) + })); + + assert_eq!(runner.drain_all(0).unwrap().abandoned, 1); + // One attempt, not MAX_ATTEMPTS, and never claimable again. + assert_eq!(runner.drain_all(1_000_000).unwrap().ran(), 0); + let (state, attempts): (i64, i64) = c + .query_row("SELECT state, attempts FROM jobs", [], |r| { + Ok((r.get(0)?, r.get(1)?)) + }) + .unwrap(); + assert_eq!(state, JobState::Failed as i64); + assert_eq!(attempts, 1); + } + + #[test] + fn an_interrupted_job_costs_no_attempt_and_ends_the_drain() { + // `onStopped()` says nothing about the file. Charging it an attempt + // would let five backgroundings mark good work as failed. + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + queued(&c, JobKind::Thumbnail, 2); + + let report = Runner::new(&c) + .with(handler(&[JobKind::Thumbnail], |_| Outcome::Interrupted)) + .drain_all(0) + .unwrap(); + + assert_eq!(report.interrupted, 1); + assert_eq!( + report.stopped, + Stopped::Cancelled, + "an interrupted job must end the drain, or releasing it hands it \ + straight back and the loop never ends" + ); + let (state, attempts): (i64, i64) = c + .query_row( + "SELECT state, attempts FROM jobs WHERE subject_id = 1", + [], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert_eq!(state, JobState::Pending as i64); + assert_eq!(attempts, 0, "the claim's speculative attempt is given back"); + } + + #[test] + fn only_kinds_a_handler_covers_are_claimed() { + // A device with no connector must leave `FetchOriginal` where it is. + // Claiming it to fail it would cost five attempts and five backoffs + // per photograph, on battery, to reach a conclusion known in advance. + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + queued(&c, JobKind::FetchOriginal, 2); + + let seen = Arc::new(Mutex::new(Vec::new())); + let report = Runner::new(&c) + .with(recording(vec![JobKind::Thumbnail], seen.clone())) + .drain_all(0) + .unwrap(); + + assert_eq!(report.completed, 1); + assert_eq!(*seen.lock().unwrap(), vec![1]); + + let (state, attempts): (i64, i64) = c + .query_row( + "SELECT state, attempts FROM jobs WHERE subject_id = 2", + [], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert_eq!(state, JobState::Pending as i64); + assert_eq!(attempts, 0, "an unhandled job is untouched, not failed"); + } + + #[test] + fn a_runner_with_no_handlers_claims_nothing() { + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + + assert_eq!(Runner::new(&c).drain_all(0).unwrap().ran(), 0); + assert_eq!(rows(&c), 1); + } + + #[test] + fn a_drain_stops_at_its_job_budget() { + let c = db(); + for id in 1..=5 { + queued(&c, JobKind::Thumbnail, id); + } + + let seen = Arc::new(Mutex::new(Vec::new())); + let mut runner = Runner::new(&c).with(recording(vec![JobKind::Thumbnail], seen.clone())); + let report = runner + .drain(&|| 0, Budget::jobs(2), &AtomicBool::new(false)) + .unwrap(); + + assert_eq!(report.completed, 2); + assert_eq!(report.stopped, Stopped::Budget); + assert!(report.stopped.more_to_do()); + assert_eq!(rows(&c), 3, "the rest is still queued for the next slot"); + assert_eq!(seen.lock().unwrap().len(), 2); + } + + #[test] + fn a_drain_stops_at_its_deadline() { + // What a `WorkManager` slot does: a fixed window, and whatever did not + // fit stays queued for the next one. + let c = db(); + for id in 1..=5 { + queued(&c, JobKind::Thumbnail, id); + } + + // A clock that advances a second per reading, so the deadline arrives + // without the test sleeping. + let tick = std::cell::Cell::new(0i64); + let clock = || { + let t = tick.get(); + tick.set(t + 1); + t + }; + + let report = Runner::new(&c) + .with(handler(&[JobKind::Thumbnail], |_| Outcome::Done)) + .drain(&clock, Budget::until(3), &AtomicBool::new(false)) + .unwrap(); + + assert_eq!(report.stopped, Stopped::Deadline); + assert_eq!(report.completed, 3, "one job per second up to the deadline"); + assert_eq!(rows(&c), 2); + } + + #[test] + fn a_cancelled_drain_stops_between_jobs() { + let c = db(); + for id in 1..=5 { + queued(&c, JobKind::Thumbnail, id); + } + + let cancel = AtomicBool::new(false); + // Cancelled from inside the handler, standing in for the host thread + // setting the flag while a job is in flight: the job in hand finishes, + // and nothing further is claimed. + let report = Runner::new(&c) + .with(handler(&[JobKind::Thumbnail], |_| { + cancel.store(true, Ordering::Relaxed); + Outcome::Done + })) + .drain(&|| 0, Budget::UNLIMITED, &cancel) + .unwrap(); + + assert_eq!(report.completed, 1); + assert_eq!(report.stopped, Stopped::Cancelled); + assert_eq!(rows(&c), 4, "the work is kept, not lost"); + } + + #[test] + fn a_job_interrupted_by_process_death_is_reclaimed_and_run_once() { + // FR-PLAT-AND-3. The kill happens between the claim and the outcome, + // which is the window a durable queue exists to survive: no `complete`, + // no `fail`, just a row marked `Running` with nobody holding it. + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + + // The dead process. It claimed the job and never came back. + let claimed = jobs::claim_next(&c, 0).unwrap().expect("claimable"); + assert_eq!(claimed.subject_id, Some(1)); + + // A fresh runner, before it starts, finds the queue empty — the row is + // `Running` and no claim will touch it. + let seen = Arc::new(Mutex::new(Vec::new())); + let mut runner = Runner::new(&c).with(recording(vec![JobKind::Thumbnail], seen.clone())); + assert_eq!( + runner.drain_all(0).unwrap().ran(), + 0, + "an orphan is invisible until it is recovered — which is exactly \ + why recovery has to happen at startup" + ); + + let recovered = runner.recover().unwrap(); + assert_eq!(recovered.reclaimed, 1); + + let report = runner.drain_all(0).unwrap(); + assert_eq!(report.completed, 1); + assert_eq!( + *seen.lock().unwrap(), + vec![1], + "resumed, not repeated and not lost" + ); + assert_eq!(rows(&c), 0); + } + + #[test] + fn a_crash_still_costs_an_attempt() { + // Deliberate: a job that takes the process down with it every time is + // indistinguishable from one that fails, and the attempt counter is + // the only evidence we keep across a death. Without this a poison-pill + // job would be reclaimed and re-run forever. + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + + for _ in 0..MAX_ATTEMPTS { + jobs::claim_next(&c, 0).unwrap().expect("claimable"); + recover(&c).unwrap(); + } + + let job = jobs::claim_next(&c, 0).unwrap().unwrap(); + assert!(job.attempts > MAX_ATTEMPTS); + jobs::fail(&c, &job, 0, "died again").unwrap(); + let state: i64 = c + .query_row("SELECT state FROM jobs", [], |r| r.get(0)) + .unwrap(); + assert_eq!(state, JobState::Failed as i64); + } + + #[test] + fn recovery_drops_jobs_whose_photograph_is_gone() { + // A culled library leaves thumbnail jobs for images deleted months + // ago. Every one would be claimed, run and failed. + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + queued(&c, JobKind::Thumbnail, 2); + c.execute("DELETE FROM images WHERE id = 2", []).unwrap(); + + let recovered = recover(&c).unwrap(); + assert_eq!(recovered.reaped, 1); + assert!(recovered.did_anything()); + + let left: i64 = c + .query_row("SELECT subject_id FROM jobs", [], |r| r.get(0)) + .unwrap(); + assert_eq!(left, 1, "only the job whose subject survives is kept"); + } + + #[test] + fn a_quiet_startup_recovers_nothing() { + let c = db(); + queued(&c, JobKind::Thumbnail, 1); + assert_eq!(recover(&c).unwrap(), Recovered::default()); + assert!(!recover(&c).unwrap().did_anything()); + } + + #[test] + fn two_runners_on_one_catalog_never_take_the_same_job() { + // Sequential rather than threaded, so the property is asserted without + // depending on the scheduler: whatever the second connection claims, + // it is not what the first one is holding. + let path = temp_catalog("contention-pair"); + let a = open(&path); + let b = open(&path); + + for id in 1..=2 { + queued(&a, JobKind::Thumbnail, id); + } + + let first = jobs::claim_next(&a, 0).unwrap().expect("one for A"); + let second = jobs::claim_next(&b, 0).unwrap().expect("one for B"); + + assert_ne!(first.id, second.id); + assert!( + jobs::claim_next(&a, 0).unwrap().is_none(), + "a claimed job is invisible to every connection, not just its own" + ); + } + + #[test] + fn concurrent_runners_share_the_queue_without_repeating_work() { + // The claim is one atomic statement precisely so this holds: four + // threads, four connections, and every job run exactly once. + const THREADS: usize = 4; + const JOBS: i64 = 24; + + let path = temp_catalog("contention-threads"); + let seeder = open(&path); + for id in 1..=JOBS { + queued(&seeder, JobKind::Thumbnail, id); + } + drop(seeder); + + let seen = Arc::new(Mutex::new(Vec::new())); + let mut threads = Vec::new(); + for _ in 0..THREADS { + let path = path.clone(); + let seen = seen.clone(); + threads.push(std::thread::spawn(move || { + let conn = open(&path); + Runner::new(&conn) + .with(recording(vec![JobKind::Thumbnail], seen)) + .drain_all(0) + .unwrap() + .completed + })); + } + + let completed: usize = threads.into_iter().map(|t| t.join().unwrap()).sum(); + assert_eq!(completed, JOBS as usize); + + let mut ran = seen.lock().unwrap().clone(); + ran.sort_unstable(); + assert_eq!( + ran, + (1..=JOBS).collect::>(), + "every job exactly once — no duplicate claim, nothing dropped" + ); + + let leftover = open(&path); + assert_eq!(rows(&leftover), 0); + } + + #[test] + fn priority_survives_the_runner() { + // NFR-ARCH-2: visible work strictly preempts bulk work, and it has to + // still be true when the queue is drained through a handler rather + // than by hand. + let c = db(); + image(&c, 1); + image(&c, 2); + enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap(); + enqueue(&c, JobKind::Thumbnail, Some(2), Priority::Interactive, None).unwrap(); + + let seen = Arc::new(Mutex::new(Vec::new())); + Runner::new(&c) + .with(recording(vec![JobKind::Thumbnail], seen.clone())) + .drain_all(0) + .unwrap(); + + assert_eq!(*seen.lock().unwrap(), vec![2, 1]); + } + + #[test] + fn a_handler_sees_the_payload_and_the_attempt_count() { + // Both are how a handler decides what to do: the payload is the only + // thing that survives from the enqueue site, and the attempt count is + // how it can tell a first try from a last one. + let c = db(); + image(&c, 1); + enqueue( + &c, + JobKind::ScanFolder, + Some(1), + Priority::Background, + Some("/lib/2024"), + ) + .unwrap(); + + let payload = Arc::new(Mutex::new(None)); + let recorded = payload.clone(); + Runner::new(&c) + .with(handler(&[JobKind::ScanFolder], move |job| { + *recorded.lock().unwrap() = Some((job.payload.clone(), job.attempts)); + Outcome::Done + })) + .drain_all(0) + .unwrap(); + + assert_eq!( + *payload.lock().unwrap(), + Some((Some("/lib/2024".to_string()), 1)) + ); + } +} diff --git a/ui/dr-ui/src/library_ui.rs b/ui/dr-ui/src/library_ui.rs index 15224f2..3be00de 100644 --- a/ui/dr-ui/src/library_ui.rs +++ b/ui/dr-ui/src/library_ui.rs @@ -1936,6 +1936,29 @@ fn show_catalog_now( } }; + // Before anything else can see this catalog, and exactly once per open — + // the early return above is what makes it once. The job queue is durable, + // so a run that was killed mid-job left its row marked `Running` with + // nobody holding it; recovery hands those back to be resumed rather than + // lost, and drops jobs naming photographs that have since been deleted. + // + // Here rather than wherever a runner starts, because there is no owner + // column in `jobs`: a second recovery pass while a worker held a claim + // would take that claim away from it. + match dr_catalog::runner::recover(cat.connection()) { + Ok(r) if r.did_anything() => log::info!( + "job queue: {} interrupted job(s) resumed, \ + {} for deleted photographs dropped", + r.reclaimed, + r.reaped + ), + Ok(_) => {} + // Not surfaced. The queue is rebuildable like everything else in the + // catalog, and a grid that refuses to open because a background queue + // could not be tidied is the worse failure by some way. + Err(e) => log::warn!("job queue could not be recovered: {e}"), + } + // The sidebar before the grid, because the grid's badges read collection // membership — the same order the scan's completion uses. crate::collections_ui::refresh_tree(window, coll_ctl, &cat);