Files
DarkRoom/core/dr-catalog/src/runner.rs
T
dtourolleandClaude Opus 5 758436cc28 Keep the runner's borrow alive as long as the connection it reads
The first compile this branch had. One borrow error, in the four-thread
contention test: the `Runner` was the block's tail expression, and a
tail's temporaries are dropped after the block's locals, so it outlived
the `conn` it borrowed. Bound to a local, with the ordering rule written
down beside it -- it is exactly the shape someone tidies back.

Everything else stood: clippy clean at -D warnings, and all 18 runner
tests pass, including the four-thread four-connection claim and the
`UPDATE ... RETURNING` rewrite the author flagged as the riskiest line
in the diff.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-29 23:48:15 +02:00

936 lines
34 KiB
Rust

//! 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<usize>,
/// Stop once the clock reaches this second. Same clock the drain is given.
pub deadline: Option<i64>,
}
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<Recovered, CatalogError> {
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<Box<dyn JobHandler + 'a>>,
/// 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<JobKind>,
}
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<dyn JobHandler + 'a>) -> 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<Recovered, CatalogError> {
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<Option<Ran>, 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<DrainReport, CatalogError> {
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<DrainReport, CatalogError> {
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<F> {
kinds: Vec<JobKind>,
act: F,
}
impl<F: FnMut(&Job) -> Outcome> JobHandler for Fake<F> {
fn kinds(&self) -> &[JobKind] {
&self.kinds
}
fn run(&mut self, _conn: &Connection, job: &Job) -> Outcome {
(self.act)(job)
}
}
fn handler<F: FnMut(&Job) -> Outcome>(kinds: &[JobKind], act: F) -> Fake<F> {
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<JobKind>, seen: Arc<Mutex<Vec<i64>>>) -> 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::<Vec<_>>(),
"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))
);
}
}