//! 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); // Bound to a local rather than left as the block's tail: the // `Runner` borrows `conn`, and a tail expression's temporaries // are dropped *after* the block's locals, so the borrow would // outlive what it borrows. let completed = Runner::new(&conn) .with(recording(vec![JobKind::Thumbnail], seen)) .drain_all(0) .unwrap() .completed; 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)) ); } }