Stopping the enqueue leaves the rows already queued: 23,582 on the reference catalog, about 1 MB of table and indexes that every query over jobs pays for. A migration would be the usual tool and is the wrong one here. A schema bump makes an older build refuse the synced catalog snapshot, and the tablet is on 0.16.0. So the rows are dropped at runtime instead, by jobs::drop_retired over a new JobKind::RETIRED list, from runner::recover - which already runs exactly once per catalog open, before any worker. It runs every open rather than once because an older build sharing the catalog queues them again on its next scan. kind leads the UNIQUE(kind, subject_id) index, so with nothing left it is one index probe. Measured on a copy of the reference catalog: 23,582 rows dropped in 40 ms on the first open, 0.07 ms after. Thumbnail stays in the enum so its number is never reused for a kind that would then inherit old rows. The runner tests that call recover move to a live kind; the jobs.rs tests of queue mechanics never call it and are unchanged. Refs #73
982 lines
36 KiB
Rust
982 lines
36 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 — face detection
|
|
/// 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,
|
|
/// Jobs deleted because their kind is retired ([`JobKind::RETIRED`]).
|
|
pub retired: usize,
|
|
}
|
|
|
|
impl Recovered {
|
|
pub fn did_anything(&self) -> bool {
|
|
self.reclaimed > 0 || self.reaped > 0 || self.retired > 0
|
|
}
|
|
}
|
|
|
|
/// Ready the queue for a fresh run, before any worker touches it.
|
|
///
|
|
/// Three distinct cleanups, and all 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 jobs for photographs that were deleted
|
|
/// months ago, and every one of them would be claimed, run and failed.
|
|
/// - **Retire.** Rows of a kind nothing enqueues or claims any more
|
|
/// ([`jobs::drop_retired`]). Here rather than in a migration so that no
|
|
/// schema bump locks an older device out of the synced catalog, and every
|
|
/// time rather than once because an older build sharing the catalog will
|
|
/// queue them again.
|
|
///
|
|
/// Retiring runs first, so the other two never touch rows about to go.
|
|
/// Reclaim runs next 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> {
|
|
let retired = jobs::drop_retired(conn)?;
|
|
Ok(Recovered {
|
|
reclaimed: jobs::recover_orphaned(conn)?,
|
|
reaped: jobs::reap_orphan_subjects(conn)?,
|
|
retired,
|
|
})
|
|
}
|
|
|
|
/// 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::ContentHash, 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::ContentHash], 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::ContentHash, 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_retired_kinds_every_time_and_nothing_else() {
|
|
// What 0.16.0 left behind: a thumbnail job per photograph that nothing
|
|
// would ever claim, beside live work that must survive.
|
|
let c = db();
|
|
queued(&c, JobKind::Thumbnail, 1);
|
|
queued(&c, JobKind::Thumbnail, 2);
|
|
enqueue(
|
|
&c,
|
|
JobKind::DetectFaces,
|
|
Some(1),
|
|
Priority::Background,
|
|
None,
|
|
)
|
|
.unwrap();
|
|
|
|
let first = recover(&c).unwrap();
|
|
assert_eq!(first.retired, 2);
|
|
assert!(first.did_anything());
|
|
|
|
let kinds: Vec<i64> = c
|
|
.prepare("SELECT kind FROM jobs")
|
|
.unwrap()
|
|
.query_map([], |r| r.get(0))
|
|
.unwrap()
|
|
.map(Result::unwrap)
|
|
.collect();
|
|
assert_eq!(kinds, vec![JobKind::DetectFaces as i64]);
|
|
|
|
// An older build opening the same catalog queues them again on its
|
|
// next scan. The next open by this one clears them again.
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
assert_eq!(recover(&c).unwrap().retired, 1);
|
|
assert_eq!(recover(&c).unwrap(), Recovered::default());
|
|
}
|
|
|
|
#[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::ContentHash, 1);
|
|
queued(&c, JobKind::ContentHash, 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::ContentHash, 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))
|
|
);
|
|
}
|
|
}
|