diff --git a/core/dr-catalog/src/jobs.rs b/core/dr-catalog/src/jobs.rs index 6913791..2bede77 100644 --- a/core/dr-catalog/src/jobs.rs +++ b/core/dr-catalog/src/jobs.rs @@ -10,8 +10,14 @@ //! table absorb the redundancy. //! - **Priority shared with the GPU scheduler** (ARCH §5.3), so one notion of //! urgency governs the whole app and visible work always preempts bulk work. +//! +//! Nothing in this file runs a job. [`crate::runner`] is the other half — the +//! one that claims from this table, does the work through a handler, and +//! reports back. Worth knowing because for a long time it did not exist: every +//! producer called [`enqueue`] and nothing ever called [`claim_next`], so the +//! table only ever grew. -use rusqlite::Connection; +use rusqlite::{Connection, OptionalExtension}; use crate::error::CatalogError; @@ -48,6 +54,21 @@ pub enum JobKind { } impl JobKind { + /// Every kind, so code that has to enumerate them cannot quietly miss one + /// that was added later. A `match` would catch that; a hand-written array + /// at each call site would not. + pub const ALL: [JobKind; 9] = [ + JobKind::ScanFolder, + JobKind::ExtractMetadata, + JobKind::Thumbnail, + JobKind::ReadSidecar, + JobKind::WriteSidecar, + JobKind::ContentHash, + JobKind::FetchPreview, + JobKind::FetchOriginal, + JobKind::DetectFaces, + ]; + fn from_i64(v: i64) -> Option { Some(match v { 0 => JobKind::ScanFolder, @@ -68,6 +89,16 @@ impl JobKind { pub fn is_network(self) -> bool { matches!(self, JobKind::FetchPreview | JobKind::FetchOriginal) } + + /// Whether `subject_id` names a row in `images`. + /// + /// Every kind but one is per-photograph. `ScanFolder`'s subject is a + /// *folder*, and the two id spaces are unrelated — so anything that joins + /// `subject_id` against `images` has to exclude it, or it will read one + /// table's ids as another's and act on the answer. + pub fn subject_is_image(self) -> bool { + !matches!(self, JobKind::ScanFolder) + } } /// Scheduling class, matching the GPU tile scheduler (ARCH §5.3). @@ -86,6 +117,22 @@ pub enum Priority { Interactive = 2, } +impl Priority { + /// Read back from the stored column. + /// + /// An unrecognised value reads as `Background` rather than failing: a + /// priority is a hint about ordering, and refusing to run a job because + /// its urgency is spelled oddly would be a worse answer than running it + /// last. + fn from_i64(v: i64) -> Self { + match v { + 2 => Priority::Interactive, + 1 => Priority::Prefetch, + _ => Priority::Background, + } + } +} + /// Lifecycle state. #[derive(Debug, Clone, Copy, PartialEq, Eq)] #[repr(i64)] @@ -156,52 +203,115 @@ pub fn enqueue( /// Claim the next runnable job, highest priority first. /// /// `now` is passed rather than read from the clock so backoff is testable. -/// Claiming marks the row `Running` in the same transaction as the read, so -/// two workers cannot take the same job. pub fn claim_next(conn: &Connection, now: i64) -> Result, CatalogError> { - let tx = conn.unchecked_transaction()?; + claim(conn, now, None) +} - let job = tx - .query_row( - "SELECT id, kind, subject_id, priority, attempts, payload - FROM jobs - WHERE state = 0 AND not_before <= ?1 - ORDER BY priority DESC, id ASC - LIMIT 1", - [now], - |r| { - Ok(( - r.get::<_, i64>(0)?, - r.get::<_, i64>(1)?, - r.get::<_, Option>(2)?, - r.get::<_, i64>(3)?, - r.get::<_, i64>(4)?, - r.get::<_, Option>(5)?, - )) - }, - ) - .ok(); +/// Claim the next runnable job of one of `kinds`. +/// +/// What lets a runner take only the work it can actually do. A device with no +/// connector must leave `FetchOriginal` rows alone rather than claim them and +/// fail them five times each with backoff; a runner gated onto an unmetered +/// network (FR-NC-6) passes the local kinds only and leaves the transfers +/// where they are. Neither is expressible by filtering *after* a claim, +/// because the claim has already marked the row `Running`. +/// +/// An empty list claims nothing, which is the honest reading of "there is +/// nothing this worker can do". +pub fn claim_next_matching( + conn: &Connection, + now: i64, + kinds: &[JobKind], +) -> Result, CatalogError> { + if kinds.is_empty() { + return Ok(None); + } + claim(conn, now, Some(kinds)) +} - let Some((id, kind, subject_id, priority, attempts, payload)) = job else { +/// The claim, as one statement. +/// +/// # Why this is not a transaction around a read and a write +/// +/// It used to be, and under two connections that is not safe in the way it +/// looks. A deferred transaction takes a read lock for the `SELECT` and only +/// tries to upgrade at the `UPDATE`; with WAL, a second worker that read the +/// same snapshot gets `SQLITE_BUSY_SNAPSHOT` on its write — an error a busy +/// handler cannot retry away, because the fix is to roll back and start over. +/// So the queue was correct only in the sense that the loser failed loudly. +/// +/// `UPDATE ... WHERE id = (SELECT ...) RETURNING` is one statement, so it is +/// one implicit transaction that takes the write lock immediately. Two workers +/// serialise, the loser waits out its busy timeout rather than erroring, and +/// neither can see a row the other is already holding. +fn claim( + conn: &Connection, + now: i64, + kinds: Option<&[JobKind]>, +) -> Result, CatalogError> { + // `now` first, then the kinds, matching the order the placeholders appear + // in the text below. + let mut args: Vec = vec![now]; + let filter = match kinds { + None => String::new(), + Some(kinds) => { + // Built from the kind *count*, never from anything a user typed — + // the same discipline `collections::descendants` uses, since + // `carray` is not compiled in. + let placeholders = std::iter::repeat_n("?", kinds.len()) + .collect::>() + .join(","); + args.extend(kinds.iter().map(|k| *k as i64)); + format!(" AND kind IN ({placeholders})") + } + }; + + let sql = format!( + "UPDATE jobs + SET state = 1, attempts = attempts + 1 + WHERE id = (SELECT id + FROM jobs + WHERE state = 0 AND not_before <= ?{filter} + ORDER BY priority DESC, id ASC + LIMIT 1) + RETURNING id, kind, subject_id, priority, attempts, payload" + ); + + let claimed = conn + .query_row(&sql, rusqlite::params_from_iter(args.iter()), |r| { + Ok(( + r.get::<_, i64>(0)?, + r.get::<_, i64>(1)?, + r.get::<_, Option>(2)?, + r.get::<_, i64>(3)?, + r.get::<_, i64>(4)?, + r.get::<_, Option>(5)?, + )) + }) + .optional()?; + + let Some((id, kind, subject_id, priority, attempts, payload)) = claimed else { return Ok(None); }; - tx.execute( - "UPDATE jobs SET state = 1, attempts = attempts + 1 WHERE id = ?1", - [id], - )?; - tx.commit()?; + let Some(kind) = JobKind::from_i64(kind) else { + // A row written by a build that knows a kind this one does not — a + // downgrade, or a catalog synced from a newer device. Running it as + // some other kind would be worse than not running it, so it is parked + // where the next claim will not see it again. + // + // Answering `None` understates what is queued for one pass. The + // alternative is a loop that keeps claiming the same unreadable row. + abandon(conn, id, &format!("unknown job kind {kind}"))?; + return Ok(None); + }; Ok(Some(Job { id, - kind: JobKind::from_i64(kind).unwrap_or(JobKind::ExtractMetadata), + kind, subject_id, - priority: match priority { - 2 => Priority::Interactive, - 1 => Priority::Prefetch, - _ => Priority::Background, - }, - attempts: attempts + 1, + priority: Priority::from_i64(priority), + attempts, payload, })) } @@ -215,16 +325,53 @@ pub fn complete(conn: &Connection, id: i64) -> Result<(), CatalogError> { /// Job failed. Reschedules with backoff, or gives up past [`MAX_ATTEMPTS`]. pub fn fail(conn: &Connection, job: &Job, now: i64, err: &str) -> Result<(), CatalogError> { if job.attempts >= MAX_ATTEMPTS { - conn.execute( - "UPDATE jobs SET state = 2, last_error = ?2 WHERE id = ?1", - rusqlite::params![job.id, err], - )?; + abandon(conn, job.id, err) } else { conn.execute( "UPDATE jobs SET state = 0, not_before = ?2, last_error = ?3 WHERE id = ?1", rusqlite::params![job.id, now + backoff_seconds(job.attempts), err], )?; + Ok(()) } +} + +/// Give up on a job now, with no further retries. +/// +/// For failures a retry cannot fix — the subject is gone, the payload is +/// unreadable, the format is one this build does not know. Walking the whole +/// retry ladder to reach a conclusion the first attempt already reached costs +/// five wakeups and five backoffs per photograph, which on a phone is the +/// difference the user notices. +/// +/// The row is kept rather than deleted, because "this file failed and here is +/// why" is something the user is entitled to see (NFR-ARCH-4). +pub fn abandon(conn: &Connection, id: i64, err: &str) -> Result<(), CatalogError> { + conn.execute( + "UPDATE jobs SET state = 2, last_error = ?2 WHERE id = ?1", + rusqlite::params![id, err], + )?; + Ok(()) +} + +/// Put a claimed job back exactly as it was found. +/// +/// For a worker that is being stopped rather than a job that is going wrong: +/// the platform revoked the slot, the user left the screen. The attempt the +/// claim consumed is given back, because nothing was learned about the file — +/// without that, five backgroundings in a row would mark good work as failed. +/// +/// Guarded on the claim still being *this* claim. There is no owner column, so +/// `attempts` stands in for one: it is bumped by every claim, so the row only +/// still reads `state = 1` with the caller's own attempt number while nobody +/// else has taken it since. A worker that comes back after the recovery pass +/// handed its job to someone else therefore changes nothing, rather than +/// releasing a job another worker is in the middle of. +pub fn release(conn: &Connection, job: &Job) -> Result<(), CatalogError> { + conn.execute( + "UPDATE jobs SET state = 0, attempts = max(0, attempts - 1) + WHERE id = ?1 AND state = 1 AND attempts = ?2", + rusqlite::params![job.id, job.attempts], + )?; Ok(()) } @@ -232,11 +379,97 @@ pub fn fail(conn: &Connection, job: &Job, now: i64, err: &str) -> Result<(), Cat /// /// A row left `Running` has no owner — the process that claimed it is gone. /// Called at startup, before any worker begins (FR-PLAT-AND-3). +/// +/// The attempt the dead claim consumed is deliberately *not* refunded. 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 — +/// without it a poison-pill job is reclaimed and re-run forever. pub fn recover_orphaned(conn: &Connection) -> Result { let n = conn.execute("UPDATE jobs SET state = 0 WHERE state = 1", [])?; Ok(n) } +/// Delete jobs whose photograph is gone. +/// +/// Coalescing keeps the table one row per unit of work, but nothing shrinks it +/// when the work stops existing: a library that has been culled carries a +/// thumbnail job for every photograph deleted since the last time anything +/// looked. Each one would be claimed, run, and failed five times. +/// +/// Only kinds whose subject really is an image ([`JobKind::subject_is_image`]) +/// are considered — `ScanFolder`'s subject is a folder id, and joining it +/// against `images` would delete jobs by coincidence of numbering. +/// +/// There is no foreign key to do this instead. `jobs.subject_id` deliberately +/// references nothing: it means different tables for different kinds, and a +/// constraint that is right for eight of nine kinds is not a constraint. +pub fn reap_orphan_subjects(conn: &Connection) -> Result { + let kinds: Vec = JobKind::ALL + .iter() + .filter(|k| k.subject_is_image()) + .map(|k| *k as i64) + .collect(); + let placeholders = std::iter::repeat_n("?", kinds.len()) + .collect::>() + .join(","); + + let n = conn.execute( + &format!( + "DELETE FROM jobs + WHERE subject_id IS NOT NULL + AND kind IN ({placeholders}) + AND NOT EXISTS (SELECT 1 FROM images WHERE images.id = jobs.subject_id)" + ), + rusqlite::params_from_iter(kinds.iter()), + )?; + Ok(n) +} + +/// How much is left, by state. +/// +/// One query rather than a listing, because the caller is a progress line: a +/// foreground service's notification has to say how much remains without +/// reading a hundred thousand rows to find out (FR-PLAT-AND-4). +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub struct Counts { + /// Claimable now or after a backoff. + pub pending: usize, + /// Claimed by someone. After a clean start that is a live worker; before + /// [`recover_orphaned`] it is a dead one. + pub running: usize, + /// Given up on, and kept so the user can see what failed and why. + pub failed: usize, +} + +impl Counts { + /// Work that is still going to happen. + pub fn outstanding(&self) -> usize { + self.pending + self.running + } +} + +/// Count the queue by state. +pub fn counts(conn: &Connection) -> Result { + // `sum` over no rows is NULL, not 0 — an empty queue would otherwise fail + // to convert rather than counting nothing. + let (pending, running, failed) = conn.query_row( + "SELECT sum(state = 0), sum(state = 1), sum(state = 2) FROM jobs", + [], + |r| { + Ok(( + r.get::<_, Option>(0)?, + r.get::<_, Option>(1)?, + r.get::<_, Option>(2)?, + )) + }, + )?; + Ok(Counts { + pending: pending.unwrap_or(0) as usize, + running: running.unwrap_or(0) as usize, + failed: failed.unwrap_or(0) as usize, + }) +} + #[cfg(test)] mod tests { use super::*; @@ -391,6 +624,241 @@ mod tests { assert_eq!(backoff_seconds(100), 300); } + /// 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(); + } + + #[test] + fn a_runner_claims_only_the_kinds_it_names() { + // The property a filtered claim exists for: work this worker cannot do + // is left untouched — not claimed, not attempted, not failed. + let c = db(); + enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap(); + enqueue( + &c, + JobKind::FetchOriginal, + Some(1), + Priority::Interactive, + None, + ) + .unwrap(); + + // `FetchOriginal` is the higher priority and would be claimed first by + // an unfiltered claim. It is not this worker's to take. + let job = claim_next_matching(&c, 0, &[JobKind::Thumbnail]) + .unwrap() + .expect("the thumbnail is claimable"); + assert_eq!(job.kind, JobKind::Thumbnail); + assert!(claim_next_matching(&c, 0, &[JobKind::Thumbnail]) + .unwrap() + .is_none()); + + let (state, attempts): (i64, i64) = c + .query_row( + "SELECT state, attempts FROM jobs WHERE kind = ?1", + [JobKind::FetchOriginal as i64], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert_eq!(state, JobState::Pending as i64); + assert_eq!(attempts, 0); + } + + #[test] + fn a_worker_that_can_do_nothing_claims_nothing() { + let c = db(); + enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap(); + assert!(claim_next_matching(&c, 0, &[]).unwrap().is_none()); + } + + #[test] + fn releasing_a_claim_gives_the_attempt_back() { + // A stopped worker has learned nothing about the file, so the claim it + // is handing back must cost nothing. + let c = db(); + enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap(); + + let job = claim_next(&c, 0).unwrap().unwrap(); + assert_eq!(job.attempts, 1); + release(&c, &job).unwrap(); + + let again = claim_next(&c, 0).unwrap().expect("claimable again at once"); + assert_eq!(again.attempts, 1, "the release refunded the first attempt"); + } + + #[test] + fn releasing_a_job_someone_else_has_reclaimed_does_nothing() { + // `attempts` standing in for an owner column. A worker that comes back + // after the recovery pass handed its job to someone else must not + // release a claim that is no longer its to release — which would drop + // the live worker's job back into the queue to be run twice. + let c = db(); + enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap(); + + let stale = claim_next(&c, 0).unwrap().unwrap(); + recover_orphaned(&c).unwrap(); + let live = claim_next(&c, 0).unwrap().unwrap(); + assert_eq!(live.attempts, 2); + + release(&c, &stale).unwrap(); + + 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::Running as i64, + "the live claim still holds" + ); + assert_eq!(attempts, 2, "and its attempt was not refunded for it"); + } + + #[test] + fn abandoning_skips_the_whole_retry_ladder() { + let c = db(); + enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap(); + let job = claim_next(&c, 0).unwrap().unwrap(); + + abandon(&c, job.id, "not an image this build can read").unwrap(); + + assert!(claim_next(&c, 1_000_000).unwrap().is_none()); + let (state, attempts, err): (i64, i64, String) = c + .query_row("SELECT state, attempts, last_error FROM jobs", [], |r| { + Ok((r.get(0)?, r.get(1)?, r.get(2)?)) + }) + .unwrap(); + assert_eq!(state, JobState::Failed as i64); + assert_eq!(attempts, 1, "one attempt, not MAX_ATTEMPTS"); + // Kept, not deleted: the user is entitled to see what failed and why. + assert!(err.contains("this build can read")); + } + + #[test] + fn jobs_for_a_deleted_photograph_are_reaped() { + // Coalescing keeps the table one row per unit of work; nothing shrank + // it when the work stopped existing. + 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::Background, None).unwrap(); + enqueue( + &c, + JobKind::ExtractMetadata, + Some(2), + Priority::Background, + None, + ) + .unwrap(); + + c.execute("DELETE FROM images WHERE id = 2", []).unwrap(); + assert_eq!(reap_orphan_subjects(&c).unwrap(), 2); + + let left: i64 = c + .query_row("SELECT subject_id FROM jobs", [], |r| r.get(0)) + .unwrap(); + assert_eq!(left, 1); + } + + #[test] + fn a_folder_scan_is_not_reaped_by_image_ids() { + // `ScanFolder`'s subject is a folder. Joining it against `images` + // would delete it whenever the numbering happened not to collide — + // which, on a fresh library, is almost always. + let c = db(); + enqueue( + &c, + JobKind::ScanFolder, + Some(1), + Priority::Background, + Some("/lib/2024"), + ) + .unwrap(); + + assert_eq!(reap_orphan_subjects(&c).unwrap(), 0); + assert!(claim_next(&c, 0).unwrap().is_some()); + } + + #[test] + fn a_job_kind_from_a_newer_build_is_parked_rather_than_guessed_at() { + // A catalog synced from a device running a later build. Running an + // unknown kind as some arbitrary known one is worse than not running + // it, and the old code silently read every unknown kind as + // `ExtractMetadata`. + let c = db(); + c.execute( + "INSERT INTO jobs(kind, subject_id, priority, state) VALUES (99, 1, 0, 0)", + [], + ) + .unwrap(); + + assert!(claim_next(&c, 0).unwrap().is_none()); + + let (state, err): (i64, String) = c + .query_row("SELECT state, last_error FROM jobs", [], |r| { + Ok((r.get(0)?, r.get(1)?)) + }) + .unwrap(); + assert_eq!(state, JobState::Failed as i64); + assert!(err.contains("99"), "{err}"); + } + + #[test] + fn counts_say_what_is_left() { + // What a foreground service's notification is built from: a number, + // without reading a hundred thousand rows to find it. + let c = db(); + assert_eq!(counts(&c).unwrap(), Counts::default()); + + for id in 1..=3 { + enqueue(&c, JobKind::Thumbnail, Some(id), Priority::Background, None).unwrap(); + } + let job = claim_next(&c, 0).unwrap().unwrap(); + abandon(&c, job.id, "nope").unwrap(); + claim_next(&c, 0).unwrap().unwrap(); + + let n = counts(&c).unwrap(); + assert_eq!(n.pending, 1); + assert_eq!(n.running, 1); + assert_eq!(n.failed, 1); + assert_eq!(n.outstanding(), 2, "failed work is not outstanding work"); + } + + #[test] + fn every_kind_is_in_all() { + // `ALL` is what the reap builds its kind filter from, so a kind added + // to the enum and forgotten here would quietly stop being reaped. + for (i, kind) in JobKind::ALL.iter().enumerate() { + assert_eq!( + JobKind::from_i64(i as i64), + Some(*kind), + "ALL is out of step with the discriminants at {i}" + ); + } + assert!(JobKind::from_i64(JobKind::ALL.len() as i64).is_none()); + } + + #[test] + fn only_a_folder_scan_has_a_non_image_subject() { + assert!(!JobKind::ScanFolder.subject_is_image()); + for kind in JobKind::ALL.iter().filter(|k| **k != JobKind::ScanFolder) { + assert!(kind.subject_is_image(), "{kind:?}"); + } + } + #[test] fn network_jobs_are_identifiable_for_metered_gating() { // FR-NC-6: transfers respect unmetered-network and charging 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..d0d171f --- /dev/null +++ b/core/dr-catalog/src/runner.rs @@ -0,0 +1,935 @@ +//! 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)) + ); + } +} 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);