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>
936 lines
34 KiB
Rust
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))
|
|
);
|
|
}
|
|
}
|