Merge: drain the job queue that nothing was draining
FR-PLAT-AND-4's Rust half, and FR-PLAT-AND-3's resumability with it. The queue's claim_next, complete, fail and recover_orphaned had no callers outside their own tests, so the jobs table accumulated rows nothing ever ran. It also fixes a claim that was not safe across two connections: the deferred transaction took a read lock for the SELECT and only tried to upgrade at the UPDATE, so in WAL the second worker got SQLITE_BUSY_SNAPSHOT, which a busy handler cannot retry away. It never double-claimed, but the loser errored. Now one UPDATE ... RETURNING. No handler is wired, deliberately. The only enqueue site reachable in the shipping app produces remote thumbnail jobs already served by the async grid worker, and inventing a second network path blind is not worth a requirement reading as covered on the strength of plumbing. Verified: clippy -D warnings clean, 376 dr-catalog and 548 dr-ui tests, 18 runner tests including four-thread contention and crash recovery. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> # Conflicts: # docs/traceability.md # ui/dr-ui/src/library.rs # ui/dr-ui/src/settings_ui.rs
This commit is contained in:
+499
-31
@@ -10,8 +10,14 @@
|
|||||||
//! table absorb the redundancy.
|
//! table absorb the redundancy.
|
||||||
//! - **Priority shared with the GPU scheduler** (ARCH §5.3), so one notion of
|
//! - **Priority shared with the GPU scheduler** (ARCH §5.3), so one notion of
|
||||||
//! urgency governs the whole app and visible work always preempts bulk work.
|
//! urgency governs the whole app and visible work always preempts bulk work.
|
||||||
|
//!
|
||||||
|
//! Nothing in this file runs a job. [`crate::runner`] is the other half — the
|
||||||
|
//! one that claims from this table, does the work through a handler, and
|
||||||
|
//! reports back. Worth knowing because for a long time it did not exist: every
|
||||||
|
//! producer called [`enqueue`] and nothing ever called [`claim_next`], so the
|
||||||
|
//! table only ever grew.
|
||||||
|
|
||||||
use rusqlite::Connection;
|
use rusqlite::{Connection, OptionalExtension};
|
||||||
|
|
||||||
use crate::error::CatalogError;
|
use crate::error::CatalogError;
|
||||||
|
|
||||||
@@ -48,6 +54,21 @@ pub enum JobKind {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl JobKind {
|
impl JobKind {
|
||||||
|
/// Every kind, so code that has to enumerate them cannot quietly miss one
|
||||||
|
/// that was added later. A `match` would catch that; a hand-written array
|
||||||
|
/// at each call site would not.
|
||||||
|
pub const ALL: [JobKind; 9] = [
|
||||||
|
JobKind::ScanFolder,
|
||||||
|
JobKind::ExtractMetadata,
|
||||||
|
JobKind::Thumbnail,
|
||||||
|
JobKind::ReadSidecar,
|
||||||
|
JobKind::WriteSidecar,
|
||||||
|
JobKind::ContentHash,
|
||||||
|
JobKind::FetchPreview,
|
||||||
|
JobKind::FetchOriginal,
|
||||||
|
JobKind::DetectFaces,
|
||||||
|
];
|
||||||
|
|
||||||
fn from_i64(v: i64) -> Option<Self> {
|
fn from_i64(v: i64) -> Option<Self> {
|
||||||
Some(match v {
|
Some(match v {
|
||||||
0 => JobKind::ScanFolder,
|
0 => JobKind::ScanFolder,
|
||||||
@@ -68,6 +89,16 @@ impl JobKind {
|
|||||||
pub fn is_network(self) -> bool {
|
pub fn is_network(self) -> bool {
|
||||||
matches!(self, JobKind::FetchPreview | JobKind::FetchOriginal)
|
matches!(self, JobKind::FetchPreview | JobKind::FetchOriginal)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Whether `subject_id` names a row in `images`.
|
||||||
|
///
|
||||||
|
/// Every kind but one is per-photograph. `ScanFolder`'s subject is a
|
||||||
|
/// *folder*, and the two id spaces are unrelated — so anything that joins
|
||||||
|
/// `subject_id` against `images` has to exclude it, or it will read one
|
||||||
|
/// table's ids as another's and act on the answer.
|
||||||
|
pub fn subject_is_image(self) -> bool {
|
||||||
|
!matches!(self, JobKind::ScanFolder)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Scheduling class, matching the GPU tile scheduler (ARCH §5.3).
|
/// Scheduling class, matching the GPU tile scheduler (ARCH §5.3).
|
||||||
@@ -86,6 +117,22 @@ pub enum Priority {
|
|||||||
Interactive = 2,
|
Interactive = 2,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl Priority {
|
||||||
|
/// Read back from the stored column.
|
||||||
|
///
|
||||||
|
/// An unrecognised value reads as `Background` rather than failing: a
|
||||||
|
/// priority is a hint about ordering, and refusing to run a job because
|
||||||
|
/// its urgency is spelled oddly would be a worse answer than running it
|
||||||
|
/// last.
|
||||||
|
fn from_i64(v: i64) -> Self {
|
||||||
|
match v {
|
||||||
|
2 => Priority::Interactive,
|
||||||
|
1 => Priority::Prefetch,
|
||||||
|
_ => Priority::Background,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Lifecycle state.
|
/// Lifecycle state.
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
#[repr(i64)]
|
#[repr(i64)]
|
||||||
@@ -156,20 +203,82 @@ pub fn enqueue(
|
|||||||
/// Claim the next runnable job, highest priority first.
|
/// Claim the next runnable job, highest priority first.
|
||||||
///
|
///
|
||||||
/// `now` is passed rather than read from the clock so backoff is testable.
|
/// `now` is passed rather than read from the clock so backoff is testable.
|
||||||
/// Claiming marks the row `Running` in the same transaction as the read, so
|
|
||||||
/// two workers cannot take the same job.
|
|
||||||
pub fn claim_next(conn: &Connection, now: i64) -> Result<Option<Job>, CatalogError> {
|
pub fn claim_next(conn: &Connection, now: i64) -> Result<Option<Job>, CatalogError> {
|
||||||
let tx = conn.unchecked_transaction()?;
|
claim(conn, now, None)
|
||||||
|
}
|
||||||
|
|
||||||
let job = tx
|
/// Claim the next runnable job of one of `kinds`.
|
||||||
.query_row(
|
///
|
||||||
"SELECT id, kind, subject_id, priority, attempts, payload
|
/// What lets a runner take only the work it can actually do. A device with no
|
||||||
|
/// connector must leave `FetchOriginal` rows alone rather than claim them and
|
||||||
|
/// fail them five times each with backoff; a runner gated onto an unmetered
|
||||||
|
/// network (FR-NC-6) passes the local kinds only and leaves the transfers
|
||||||
|
/// where they are. Neither is expressible by filtering *after* a claim,
|
||||||
|
/// because the claim has already marked the row `Running`.
|
||||||
|
///
|
||||||
|
/// An empty list claims nothing, which is the honest reading of "there is
|
||||||
|
/// nothing this worker can do".
|
||||||
|
pub fn claim_next_matching(
|
||||||
|
conn: &Connection,
|
||||||
|
now: i64,
|
||||||
|
kinds: &[JobKind],
|
||||||
|
) -> Result<Option<Job>, CatalogError> {
|
||||||
|
if kinds.is_empty() {
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
claim(conn, now, Some(kinds))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The claim, as one statement.
|
||||||
|
///
|
||||||
|
/// # Why this is not a transaction around a read and a write
|
||||||
|
///
|
||||||
|
/// It used to be, and under two connections that is not safe in the way it
|
||||||
|
/// looks. A deferred transaction takes a read lock for the `SELECT` and only
|
||||||
|
/// tries to upgrade at the `UPDATE`; with WAL, a second worker that read the
|
||||||
|
/// same snapshot gets `SQLITE_BUSY_SNAPSHOT` on its write — an error a busy
|
||||||
|
/// handler cannot retry away, because the fix is to roll back and start over.
|
||||||
|
/// So the queue was correct only in the sense that the loser failed loudly.
|
||||||
|
///
|
||||||
|
/// `UPDATE ... WHERE id = (SELECT ...) RETURNING` is one statement, so it is
|
||||||
|
/// one implicit transaction that takes the write lock immediately. Two workers
|
||||||
|
/// serialise, the loser waits out its busy timeout rather than erroring, and
|
||||||
|
/// neither can see a row the other is already holding.
|
||||||
|
fn claim(
|
||||||
|
conn: &Connection,
|
||||||
|
now: i64,
|
||||||
|
kinds: Option<&[JobKind]>,
|
||||||
|
) -> Result<Option<Job>, CatalogError> {
|
||||||
|
// `now` first, then the kinds, matching the order the placeholders appear
|
||||||
|
// in the text below.
|
||||||
|
let mut args: Vec<i64> = vec![now];
|
||||||
|
let filter = match kinds {
|
||||||
|
None => String::new(),
|
||||||
|
Some(kinds) => {
|
||||||
|
// Built from the kind *count*, never from anything a user typed —
|
||||||
|
// the same discipline `collections::descendants` uses, since
|
||||||
|
// `carray` is not compiled in.
|
||||||
|
let placeholders = std::iter::repeat_n("?", kinds.len())
|
||||||
|
.collect::<Vec<_>>()
|
||||||
|
.join(",");
|
||||||
|
args.extend(kinds.iter().map(|k| *k as i64));
|
||||||
|
format!(" AND kind IN ({placeholders})")
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
let sql = format!(
|
||||||
|
"UPDATE jobs
|
||||||
|
SET state = 1, attempts = attempts + 1
|
||||||
|
WHERE id = (SELECT id
|
||||||
FROM jobs
|
FROM jobs
|
||||||
WHERE state = 0 AND not_before <= ?1
|
WHERE state = 0 AND not_before <= ?{filter}
|
||||||
ORDER BY priority DESC, id ASC
|
ORDER BY priority DESC, id ASC
|
||||||
LIMIT 1",
|
LIMIT 1)
|
||||||
[now],
|
RETURNING id, kind, subject_id, priority, attempts, payload"
|
||||||
|r| {
|
);
|
||||||
|
|
||||||
|
let claimed = conn
|
||||||
|
.query_row(&sql, rusqlite::params_from_iter(args.iter()), |r| {
|
||||||
Ok((
|
Ok((
|
||||||
r.get::<_, i64>(0)?,
|
r.get::<_, i64>(0)?,
|
||||||
r.get::<_, i64>(1)?,
|
r.get::<_, i64>(1)?,
|
||||||
@@ -178,30 +287,31 @@ pub fn claim_next(conn: &Connection, now: i64) -> Result<Option<Job>, CatalogErr
|
|||||||
r.get::<_, i64>(4)?,
|
r.get::<_, i64>(4)?,
|
||||||
r.get::<_, Option<String>>(5)?,
|
r.get::<_, Option<String>>(5)?,
|
||||||
))
|
))
|
||||||
},
|
})
|
||||||
)
|
.optional()?;
|
||||||
.ok();
|
|
||||||
|
|
||||||
let Some((id, kind, subject_id, priority, attempts, payload)) = job else {
|
let Some((id, kind, subject_id, priority, attempts, payload)) = claimed else {
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
};
|
};
|
||||||
|
|
||||||
tx.execute(
|
let Some(kind) = JobKind::from_i64(kind) else {
|
||||||
"UPDATE jobs SET state = 1, attempts = attempts + 1 WHERE id = ?1",
|
// A row written by a build that knows a kind this one does not — a
|
||||||
[id],
|
// downgrade, or a catalog synced from a newer device. Running it as
|
||||||
)?;
|
// some other kind would be worse than not running it, so it is parked
|
||||||
tx.commit()?;
|
// where the next claim will not see it again.
|
||||||
|
//
|
||||||
|
// Answering `None` understates what is queued for one pass. The
|
||||||
|
// alternative is a loop that keeps claiming the same unreadable row.
|
||||||
|
abandon(conn, id, &format!("unknown job kind {kind}"))?;
|
||||||
|
return Ok(None);
|
||||||
|
};
|
||||||
|
|
||||||
Ok(Some(Job {
|
Ok(Some(Job {
|
||||||
id,
|
id,
|
||||||
kind: JobKind::from_i64(kind).unwrap_or(JobKind::ExtractMetadata),
|
kind,
|
||||||
subject_id,
|
subject_id,
|
||||||
priority: match priority {
|
priority: Priority::from_i64(priority),
|
||||||
2 => Priority::Interactive,
|
attempts,
|
||||||
1 => Priority::Prefetch,
|
|
||||||
_ => Priority::Background,
|
|
||||||
},
|
|
||||||
attempts: attempts + 1,
|
|
||||||
payload,
|
payload,
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
@@ -215,16 +325,53 @@ pub fn complete(conn: &Connection, id: i64) -> Result<(), CatalogError> {
|
|||||||
/// Job failed. Reschedules with backoff, or gives up past [`MAX_ATTEMPTS`].
|
/// Job failed. Reschedules with backoff, or gives up past [`MAX_ATTEMPTS`].
|
||||||
pub fn fail(conn: &Connection, job: &Job, now: i64, err: &str) -> Result<(), CatalogError> {
|
pub fn fail(conn: &Connection, job: &Job, now: i64, err: &str) -> Result<(), CatalogError> {
|
||||||
if job.attempts >= MAX_ATTEMPTS {
|
if job.attempts >= MAX_ATTEMPTS {
|
||||||
conn.execute(
|
abandon(conn, job.id, err)
|
||||||
"UPDATE jobs SET state = 2, last_error = ?2 WHERE id = ?1",
|
|
||||||
rusqlite::params![job.id, err],
|
|
||||||
)?;
|
|
||||||
} else {
|
} else {
|
||||||
conn.execute(
|
conn.execute(
|
||||||
"UPDATE jobs SET state = 0, not_before = ?2, last_error = ?3 WHERE id = ?1",
|
"UPDATE jobs SET state = 0, not_before = ?2, last_error = ?3 WHERE id = ?1",
|
||||||
rusqlite::params![job.id, now + backoff_seconds(job.attempts), err],
|
rusqlite::params![job.id, now + backoff_seconds(job.attempts), err],
|
||||||
)?;
|
)?;
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Give up on a job now, with no further retries.
|
||||||
|
///
|
||||||
|
/// For failures a retry cannot fix — the subject is gone, the payload is
|
||||||
|
/// unreadable, the format is one this build does not know. Walking the whole
|
||||||
|
/// retry ladder to reach a conclusion the first attempt already reached costs
|
||||||
|
/// five wakeups and five backoffs per photograph, which on a phone is the
|
||||||
|
/// difference the user notices.
|
||||||
|
///
|
||||||
|
/// The row is kept rather than deleted, because "this file failed and here is
|
||||||
|
/// why" is something the user is entitled to see (NFR-ARCH-4).
|
||||||
|
pub fn abandon(conn: &Connection, id: i64, err: &str) -> Result<(), CatalogError> {
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE jobs SET state = 2, last_error = ?2 WHERE id = ?1",
|
||||||
|
rusqlite::params![id, err],
|
||||||
|
)?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Put a claimed job back exactly as it was found.
|
||||||
|
///
|
||||||
|
/// For a worker that is being stopped rather than a job that is going wrong:
|
||||||
|
/// the platform revoked the slot, the user left the screen. The attempt the
|
||||||
|
/// claim consumed is given back, because nothing was learned about the file —
|
||||||
|
/// without that, five backgroundings in a row would mark good work as failed.
|
||||||
|
///
|
||||||
|
/// Guarded on the claim still being *this* claim. There is no owner column, so
|
||||||
|
/// `attempts` stands in for one: it is bumped by every claim, so the row only
|
||||||
|
/// still reads `state = 1` with the caller's own attempt number while nobody
|
||||||
|
/// else has taken it since. A worker that comes back after the recovery pass
|
||||||
|
/// handed its job to someone else therefore changes nothing, rather than
|
||||||
|
/// releasing a job another worker is in the middle of.
|
||||||
|
pub fn release(conn: &Connection, job: &Job) -> Result<(), CatalogError> {
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE jobs SET state = 0, attempts = max(0, attempts - 1)
|
||||||
|
WHERE id = ?1 AND state = 1 AND attempts = ?2",
|
||||||
|
rusqlite::params![job.id, job.attempts],
|
||||||
|
)?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -232,11 +379,97 @@ pub fn fail(conn: &Connection, job: &Job, now: i64, err: &str) -> Result<(), Cat
|
|||||||
///
|
///
|
||||||
/// A row left `Running` has no owner — the process that claimed it is gone.
|
/// A row left `Running` has no owner — the process that claimed it is gone.
|
||||||
/// Called at startup, before any worker begins (FR-PLAT-AND-3).
|
/// Called at startup, before any worker begins (FR-PLAT-AND-3).
|
||||||
|
///
|
||||||
|
/// The attempt the dead claim consumed is deliberately *not* refunded. A job
|
||||||
|
/// that takes the process down with it is indistinguishable from one that
|
||||||
|
/// fails, and the attempt counter is the only evidence that survives a death —
|
||||||
|
/// without it a poison-pill job is reclaimed and re-run forever.
|
||||||
pub fn recover_orphaned(conn: &Connection) -> Result<usize, CatalogError> {
|
pub fn recover_orphaned(conn: &Connection) -> Result<usize, CatalogError> {
|
||||||
let n = conn.execute("UPDATE jobs SET state = 0 WHERE state = 1", [])?;
|
let n = conn.execute("UPDATE jobs SET state = 0 WHERE state = 1", [])?;
|
||||||
Ok(n)
|
Ok(n)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Delete jobs whose photograph is gone.
|
||||||
|
///
|
||||||
|
/// Coalescing keeps the table one row per unit of work, but nothing shrinks it
|
||||||
|
/// when the work stops existing: a library that has been culled carries a
|
||||||
|
/// thumbnail job for every photograph deleted since the last time anything
|
||||||
|
/// looked. Each one would be claimed, run, and failed five times.
|
||||||
|
///
|
||||||
|
/// Only kinds whose subject really is an image ([`JobKind::subject_is_image`])
|
||||||
|
/// are considered — `ScanFolder`'s subject is a folder id, and joining it
|
||||||
|
/// against `images` would delete jobs by coincidence of numbering.
|
||||||
|
///
|
||||||
|
/// There is no foreign key to do this instead. `jobs.subject_id` deliberately
|
||||||
|
/// references nothing: it means different tables for different kinds, and a
|
||||||
|
/// constraint that is right for eight of nine kinds is not a constraint.
|
||||||
|
pub fn reap_orphan_subjects(conn: &Connection) -> Result<usize, CatalogError> {
|
||||||
|
let kinds: Vec<i64> = JobKind::ALL
|
||||||
|
.iter()
|
||||||
|
.filter(|k| k.subject_is_image())
|
||||||
|
.map(|k| *k as i64)
|
||||||
|
.collect();
|
||||||
|
let placeholders = std::iter::repeat_n("?", kinds.len())
|
||||||
|
.collect::<Vec<_>>()
|
||||||
|
.join(",");
|
||||||
|
|
||||||
|
let n = conn.execute(
|
||||||
|
&format!(
|
||||||
|
"DELETE FROM jobs
|
||||||
|
WHERE subject_id IS NOT NULL
|
||||||
|
AND kind IN ({placeholders})
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM images WHERE images.id = jobs.subject_id)"
|
||||||
|
),
|
||||||
|
rusqlite::params_from_iter(kinds.iter()),
|
||||||
|
)?;
|
||||||
|
Ok(n)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// How much is left, by state.
|
||||||
|
///
|
||||||
|
/// One query rather than a listing, because the caller is a progress line: a
|
||||||
|
/// foreground service's notification has to say how much remains without
|
||||||
|
/// reading a hundred thousand rows to find out (FR-PLAT-AND-4).
|
||||||
|
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub struct Counts {
|
||||||
|
/// Claimable now or after a backoff.
|
||||||
|
pub pending: usize,
|
||||||
|
/// Claimed by someone. After a clean start that is a live worker; before
|
||||||
|
/// [`recover_orphaned`] it is a dead one.
|
||||||
|
pub running: usize,
|
||||||
|
/// Given up on, and kept so the user can see what failed and why.
|
||||||
|
pub failed: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Counts {
|
||||||
|
/// Work that is still going to happen.
|
||||||
|
pub fn outstanding(&self) -> usize {
|
||||||
|
self.pending + self.running
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Count the queue by state.
|
||||||
|
pub fn counts(conn: &Connection) -> Result<Counts, CatalogError> {
|
||||||
|
// `sum` over no rows is NULL, not 0 — an empty queue would otherwise fail
|
||||||
|
// to convert rather than counting nothing.
|
||||||
|
let (pending, running, failed) = conn.query_row(
|
||||||
|
"SELECT sum(state = 0), sum(state = 1), sum(state = 2) FROM jobs",
|
||||||
|
[],
|
||||||
|
|r| {
|
||||||
|
Ok((
|
||||||
|
r.get::<_, Option<i64>>(0)?,
|
||||||
|
r.get::<_, Option<i64>>(1)?,
|
||||||
|
r.get::<_, Option<i64>>(2)?,
|
||||||
|
))
|
||||||
|
},
|
||||||
|
)?;
|
||||||
|
Ok(Counts {
|
||||||
|
pending: pending.unwrap_or(0) as usize,
|
||||||
|
running: running.unwrap_or(0) as usize,
|
||||||
|
failed: failed.unwrap_or(0) as usize,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
@@ -391,6 +624,241 @@ mod tests {
|
|||||||
assert_eq!(backoff_seconds(100), 300);
|
assert_eq!(backoff_seconds(100), 300);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// An image row, so a job has a subject that exists.
|
||||||
|
fn image(c: &Connection, id: i64) {
|
||||||
|
c.execute(
|
||||||
|
"INSERT INTO roots(id, kind, label) VALUES (1, 'local', '/lib')
|
||||||
|
ON CONFLICT DO NOTHING",
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
c.execute(
|
||||||
|
"INSERT INTO images(id, root_id, source_ref, added_at) VALUES (?1, 1, ?2, 0)",
|
||||||
|
rusqlite::params![id, format!("/lib/{id}.CR3")],
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_runner_claims_only_the_kinds_it_names() {
|
||||||
|
// The property a filtered claim exists for: work this worker cannot do
|
||||||
|
// is left untouched — not claimed, not attempted, not failed.
|
||||||
|
let c = db();
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
||||||
|
enqueue(
|
||||||
|
&c,
|
||||||
|
JobKind::FetchOriginal,
|
||||||
|
Some(1),
|
||||||
|
Priority::Interactive,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
// `FetchOriginal` is the higher priority and would be claimed first by
|
||||||
|
// an unfiltered claim. It is not this worker's to take.
|
||||||
|
let job = claim_next_matching(&c, 0, &[JobKind::Thumbnail])
|
||||||
|
.unwrap()
|
||||||
|
.expect("the thumbnail is claimable");
|
||||||
|
assert_eq!(job.kind, JobKind::Thumbnail);
|
||||||
|
assert!(claim_next_matching(&c, 0, &[JobKind::Thumbnail])
|
||||||
|
.unwrap()
|
||||||
|
.is_none());
|
||||||
|
|
||||||
|
let (state, attempts): (i64, i64) = c
|
||||||
|
.query_row(
|
||||||
|
"SELECT state, attempts FROM jobs WHERE kind = ?1",
|
||||||
|
[JobKind::FetchOriginal as i64],
|
||||||
|
|r| Ok((r.get(0)?, r.get(1)?)),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(state, JobState::Pending as i64);
|
||||||
|
assert_eq!(attempts, 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_worker_that_can_do_nothing_claims_nothing() {
|
||||||
|
let c = db();
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
||||||
|
assert!(claim_next_matching(&c, 0, &[]).unwrap().is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn releasing_a_claim_gives_the_attempt_back() {
|
||||||
|
// A stopped worker has learned nothing about the file, so the claim it
|
||||||
|
// is handing back must cost nothing.
|
||||||
|
let c = db();
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
||||||
|
|
||||||
|
let job = claim_next(&c, 0).unwrap().unwrap();
|
||||||
|
assert_eq!(job.attempts, 1);
|
||||||
|
release(&c, &job).unwrap();
|
||||||
|
|
||||||
|
let again = claim_next(&c, 0).unwrap().expect("claimable again at once");
|
||||||
|
assert_eq!(again.attempts, 1, "the release refunded the first attempt");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn releasing_a_job_someone_else_has_reclaimed_does_nothing() {
|
||||||
|
// `attempts` standing in for an owner column. A worker that comes back
|
||||||
|
// after the recovery pass handed its job to someone else must not
|
||||||
|
// release a claim that is no longer its to release — which would drop
|
||||||
|
// the live worker's job back into the queue to be run twice.
|
||||||
|
let c = db();
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
||||||
|
|
||||||
|
let stale = claim_next(&c, 0).unwrap().unwrap();
|
||||||
|
recover_orphaned(&c).unwrap();
|
||||||
|
let live = claim_next(&c, 0).unwrap().unwrap();
|
||||||
|
assert_eq!(live.attempts, 2);
|
||||||
|
|
||||||
|
release(&c, &stale).unwrap();
|
||||||
|
|
||||||
|
let (state, attempts): (i64, i64) = c
|
||||||
|
.query_row("SELECT state, attempts FROM jobs", [], |r| {
|
||||||
|
Ok((r.get(0)?, r.get(1)?))
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(
|
||||||
|
state,
|
||||||
|
JobState::Running as i64,
|
||||||
|
"the live claim still holds"
|
||||||
|
);
|
||||||
|
assert_eq!(attempts, 2, "and its attempt was not refunded for it");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn abandoning_skips_the_whole_retry_ladder() {
|
||||||
|
let c = db();
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
||||||
|
let job = claim_next(&c, 0).unwrap().unwrap();
|
||||||
|
|
||||||
|
abandon(&c, job.id, "not an image this build can read").unwrap();
|
||||||
|
|
||||||
|
assert!(claim_next(&c, 1_000_000).unwrap().is_none());
|
||||||
|
let (state, attempts, err): (i64, i64, String) = c
|
||||||
|
.query_row("SELECT state, attempts, last_error FROM jobs", [], |r| {
|
||||||
|
Ok((r.get(0)?, r.get(1)?, r.get(2)?))
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(state, JobState::Failed as i64);
|
||||||
|
assert_eq!(attempts, 1, "one attempt, not MAX_ATTEMPTS");
|
||||||
|
// Kept, not deleted: the user is entitled to see what failed and why.
|
||||||
|
assert!(err.contains("this build can read"));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn jobs_for_a_deleted_photograph_are_reaped() {
|
||||||
|
// Coalescing keeps the table one row per unit of work; nothing shrank
|
||||||
|
// it when the work stopped existing.
|
||||||
|
let c = db();
|
||||||
|
image(&c, 1);
|
||||||
|
image(&c, 2);
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(2), Priority::Background, None).unwrap();
|
||||||
|
enqueue(
|
||||||
|
&c,
|
||||||
|
JobKind::ExtractMetadata,
|
||||||
|
Some(2),
|
||||||
|
Priority::Background,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
c.execute("DELETE FROM images WHERE id = 2", []).unwrap();
|
||||||
|
assert_eq!(reap_orphan_subjects(&c).unwrap(), 2);
|
||||||
|
|
||||||
|
let left: i64 = c
|
||||||
|
.query_row("SELECT subject_id FROM jobs", [], |r| r.get(0))
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(left, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_folder_scan_is_not_reaped_by_image_ids() {
|
||||||
|
// `ScanFolder`'s subject is a folder. Joining it against `images`
|
||||||
|
// would delete it whenever the numbering happened not to collide —
|
||||||
|
// which, on a fresh library, is almost always.
|
||||||
|
let c = db();
|
||||||
|
enqueue(
|
||||||
|
&c,
|
||||||
|
JobKind::ScanFolder,
|
||||||
|
Some(1),
|
||||||
|
Priority::Background,
|
||||||
|
Some("/lib/2024"),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert_eq!(reap_orphan_subjects(&c).unwrap(), 0);
|
||||||
|
assert!(claim_next(&c, 0).unwrap().is_some());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_job_kind_from_a_newer_build_is_parked_rather_than_guessed_at() {
|
||||||
|
// A catalog synced from a device running a later build. Running an
|
||||||
|
// unknown kind as some arbitrary known one is worse than not running
|
||||||
|
// it, and the old code silently read every unknown kind as
|
||||||
|
// `ExtractMetadata`.
|
||||||
|
let c = db();
|
||||||
|
c.execute(
|
||||||
|
"INSERT INTO jobs(kind, subject_id, priority, state) VALUES (99, 1, 0, 0)",
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert!(claim_next(&c, 0).unwrap().is_none());
|
||||||
|
|
||||||
|
let (state, err): (i64, String) = c
|
||||||
|
.query_row("SELECT state, last_error FROM jobs", [], |r| {
|
||||||
|
Ok((r.get(0)?, r.get(1)?))
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(state, JobState::Failed as i64);
|
||||||
|
assert!(err.contains("99"), "{err}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn counts_say_what_is_left() {
|
||||||
|
// What a foreground service's notification is built from: a number,
|
||||||
|
// without reading a hundred thousand rows to find it.
|
||||||
|
let c = db();
|
||||||
|
assert_eq!(counts(&c).unwrap(), Counts::default());
|
||||||
|
|
||||||
|
for id in 1..=3 {
|
||||||
|
enqueue(&c, JobKind::Thumbnail, Some(id), Priority::Background, None).unwrap();
|
||||||
|
}
|
||||||
|
let job = claim_next(&c, 0).unwrap().unwrap();
|
||||||
|
abandon(&c, job.id, "nope").unwrap();
|
||||||
|
claim_next(&c, 0).unwrap().unwrap();
|
||||||
|
|
||||||
|
let n = counts(&c).unwrap();
|
||||||
|
assert_eq!(n.pending, 1);
|
||||||
|
assert_eq!(n.running, 1);
|
||||||
|
assert_eq!(n.failed, 1);
|
||||||
|
assert_eq!(n.outstanding(), 2, "failed work is not outstanding work");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn every_kind_is_in_all() {
|
||||||
|
// `ALL` is what the reap builds its kind filter from, so a kind added
|
||||||
|
// to the enum and forgotten here would quietly stop being reaped.
|
||||||
|
for (i, kind) in JobKind::ALL.iter().enumerate() {
|
||||||
|
assert_eq!(
|
||||||
|
JobKind::from_i64(i as i64),
|
||||||
|
Some(*kind),
|
||||||
|
"ALL is out of step with the discriminants at {i}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
assert!(JobKind::from_i64(JobKind::ALL.len() as i64).is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn only_a_folder_scan_has_a_non_image_subject() {
|
||||||
|
assert!(!JobKind::ScanFolder.subject_is_image());
|
||||||
|
for kind in JobKind::ALL.iter().filter(|k| **k != JobKind::ScanFolder) {
|
||||||
|
assert!(kind.subject_is_image(), "{kind:?}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn network_jobs_are_identifiable_for_metered_gating() {
|
fn network_jobs_are_identifiable_for_metered_gating() {
|
||||||
// FR-NC-6: transfers respect unmetered-network and charging
|
// FR-NC-6: transfers respect unmetered-network and charging
|
||||||
|
|||||||
@@ -18,6 +18,7 @@
|
|||||||
//! - [`faces`] — detected faces, the people they belong to, and who said so
|
//! - [`faces`] — detected faces, the people they belong to, and who said so
|
||||||
//! - [`bursts`] — frames that are one moment, grouped so they judge as one
|
//! - [`bursts`] — frames that are one moment, grouped so they judge as one
|
||||||
//! - [`jobs`] — the durable background work queue
|
//! - [`jobs`] — the durable background work queue
|
||||||
|
//! - [`runner`] — the thing that drains it, driven by whoever owns the thread
|
||||||
//! - [`trash`] — soft delete to a folder, then permanent delete
|
//! - [`trash`] — soft delete to a folder, then permanent delete
|
||||||
//! - [`merge`] / [`sync`] — cross-device merging of collections and keywords
|
//! - [`merge`] / [`sync`] — cross-device merging of collections and keywords
|
||||||
//!
|
//!
|
||||||
@@ -46,6 +47,7 @@ pub mod keywords;
|
|||||||
pub mod merge;
|
pub mod merge;
|
||||||
pub mod query;
|
pub mod query;
|
||||||
pub mod rating;
|
pub mod rating;
|
||||||
|
pub mod runner;
|
||||||
pub mod scan;
|
pub mod scan;
|
||||||
pub mod schema;
|
pub mod schema;
|
||||||
pub mod sync;
|
pub mod sync;
|
||||||
@@ -63,6 +65,10 @@ pub use keywords::{Coverage, Keyword, KeywordId, SelectionKeyword};
|
|||||||
pub use merge::MergeReport;
|
pub use merge::MergeReport;
|
||||||
pub use query::{Query, Sort};
|
pub use query::{Query, Sort};
|
||||||
pub use rating::{Judgement, MAX_RATING};
|
pub use rating::{Judgement, MAX_RATING};
|
||||||
|
// Not `runner::Budget`: `cache::Budget` already owns that name here and
|
||||||
|
// means something else entirely (bytes on disk, not jobs in a slot).
|
||||||
|
// Callers spell the work budget `runner::Budget`, where it is unambiguous.
|
||||||
|
pub use runner::{DrainReport, JobHandler, Outcome, Runner};
|
||||||
pub use scan::{DirAction, DirState, EntryAction, ScanOutcome};
|
pub use scan::{DirAction, DirState, EntryAction, ScanOutcome};
|
||||||
pub use trash::{TrashedImage, TRASH_DIR};
|
pub use trash::{TrashedImage, TRASH_DIR};
|
||||||
pub use walk::{ensure_root, mark_root_offline, scan_root, RootKind, ScanProgress, ScanReport};
|
pub use walk::{ensure_root, mark_root_offline, scan_root, RootKind, ScanProgress, ScanReport};
|
||||||
|
|||||||
@@ -0,0 +1,935 @@
|
|||||||
|
//! TRACES: FR-PLAT-AND-4 | FR-PLAT-AND-3
|
||||||
|
//! The thing that drains the queue.
|
||||||
|
//!
|
||||||
|
//! [`crate::jobs`] has been a complete, durable, coalescing work queue since
|
||||||
|
//! the catalog was written, and nothing has ever taken a job out of it. Every
|
||||||
|
//! producer — the local walk, the remote scan — called `enqueue` and no one
|
||||||
|
//! called `claim_next`, so the table grew one row per photograph and stayed
|
||||||
|
//! that size forever. This module is the missing half.
|
||||||
|
//!
|
||||||
|
//! # Why the runner is driven rather than self-owning
|
||||||
|
//!
|
||||||
|
//! The obvious shape is a thread that loops until the queue is empty, and it
|
||||||
|
//! is the wrong one. On Android the process does not decide when background
|
||||||
|
//! work may run: `WorkManager` does, subject to Doze, battery saver and the
|
||||||
|
//! metered-network constraints in FR-NC-6, and it revokes permission mid-job
|
||||||
|
//! by calling `onStopped()` (FR-PLAT-AND-4). A foreground service for a
|
||||||
|
//! user-initiated export gets a longer leash but still not an unbounded one.
|
||||||
|
//!
|
||||||
|
//! So the runner owns no thread, no clock and no policy. It exposes
|
||||||
|
//! [`Runner::run_one`] — claim one job, run it, record what happened — and
|
||||||
|
//! [`Runner::drain`], which repeats that against a [`Budget`] and a
|
||||||
|
//! cancellation flag the host owns. A `Worker.doWork()` that must return
|
||||||
|
//! within ten minutes calls `drain` with a deadline; a desktop idle loop calls
|
||||||
|
//! it with none. Neither has to reach inside.
|
||||||
|
//!
|
||||||
|
//! Everything the host supplies is passed in for the same reason `jobs` takes
|
||||||
|
//! `now` rather than reading the clock: a scheduler is exactly the thing that
|
||||||
|
//! has to be testable without waiting.
|
||||||
|
//!
|
||||||
|
//! # Why interruption is not failure
|
||||||
|
//!
|
||||||
|
//! Four things can happen to a claimed job, and only two of them are the job's
|
||||||
|
//! fault:
|
||||||
|
//!
|
||||||
|
//! - [`Outcome::Done`] — the row is deleted.
|
||||||
|
//! - [`Outcome::Retry`] — the work failed and might succeed later. Backoff,
|
||||||
|
//! and eventually [`crate::jobs::MAX_ATTEMPTS`] gives up on it.
|
||||||
|
//! - [`Outcome::Abandon`] — the work cannot succeed, ever. Failing five times
|
||||||
|
//! over five minutes to learn that is five minutes of a phone's battery.
|
||||||
|
//! - [`Outcome::Interrupted`] — the *host* stopped, not the job. The claim is
|
||||||
|
//! released and the attempt it consumed is given back, because a user who
|
||||||
|
//! pulled the app off the screen has not told us anything about the file.
|
||||||
|
//!
|
||||||
|
//! Process death is the fifth case and the one that cannot report itself: the
|
||||||
|
//! row simply stays `Running` with no owner. [`Runner::recover`] is what
|
||||||
|
//! reclaims it, and it is why an interrupted job is resumable rather than lost
|
||||||
|
//! (FR-PLAT-AND-3). It must run **before** any worker starts against a
|
||||||
|
//! catalog, or it will steal a job another runner is holding — there is no
|
||||||
|
//! owner column to tell them apart.
|
||||||
|
|
||||||
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
|
|
||||||
|
use rusqlite::Connection;
|
||||||
|
|
||||||
|
use crate::error::CatalogError;
|
||||||
|
use crate::jobs::{self, Job, JobKind};
|
||||||
|
|
||||||
|
/// What running a job turned out to be.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub enum Outcome {
|
||||||
|
/// The work is done. The row goes away.
|
||||||
|
Done,
|
||||||
|
/// It failed, and trying again later is reasonable. Backoff applies, and
|
||||||
|
/// [`crate::jobs::MAX_ATTEMPTS`] eventually stops it.
|
||||||
|
Retry(String),
|
||||||
|
/// It failed in a way no retry can fix — the subject is gone, the payload
|
||||||
|
/// is unreadable, the format is one this build does not know. Marked
|
||||||
|
/// failed at once rather than burning the whole retry ladder to reach the
|
||||||
|
/// same answer.
|
||||||
|
Abandon(String),
|
||||||
|
/// The host is stopping, and the job never really ran.
|
||||||
|
///
|
||||||
|
/// Distinct from `Retry` because it costs no attempt: `onStopped()` five
|
||||||
|
/// times in a row would otherwise mark a perfectly good job as failed.
|
||||||
|
Interrupted,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Something that can actually do the work a job describes.
|
||||||
|
///
|
||||||
|
/// The catalog knows what needs doing and nothing about how — a thumbnail
|
||||||
|
/// needs a decoder, a fetch needs a network stack, and neither belongs under
|
||||||
|
/// `core/dr-catalog` (ARCH §4.1: calls go downward). So the queue lives here
|
||||||
|
/// and the handlers are supplied from above.
|
||||||
|
pub trait JobHandler {
|
||||||
|
/// The kinds this handler will accept.
|
||||||
|
///
|
||||||
|
/// Load-bearing, not documentation: the runner claims **only** kinds some
|
||||||
|
/// handler declares. A queue holding `FetchOriginal` rows on a device with
|
||||||
|
/// no connector must leave them alone rather than claim them and fail
|
||||||
|
/// them, and a runner that claimed everything would do exactly that — five
|
||||||
|
/// times each, with backoff, on battery.
|
||||||
|
fn kinds(&self) -> &[JobKind];
|
||||||
|
|
||||||
|
/// Do the work.
|
||||||
|
///
|
||||||
|
/// The connection is offered because most handlers write their result back
|
||||||
|
/// into the catalog; one that does not is free to ignore it. It is the
|
||||||
|
/// runner's own connection, so a handler must not hold a transaction open
|
||||||
|
/// across a network call — the runner needs it back to record the outcome.
|
||||||
|
fn run(&mut self, conn: &Connection, job: &Job) -> Outcome;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// How much work a host is willing to let one drain do.
|
||||||
|
///
|
||||||
|
/// Both limits are checked *before* a job is claimed, never during one: a
|
||||||
|
/// handler is opaque and may be halfway through writing a sidecar. Overrunning
|
||||||
|
/// a deadline by one job is survivable; being killed mid-write is the thing
|
||||||
|
/// [`Runner::recover`] exists to clean up after.
|
||||||
|
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub struct Budget {
|
||||||
|
/// Stop after this many jobs. `None` means "until the queue is empty".
|
||||||
|
pub max_jobs: Option<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))
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1936,6 +1936,29 @@ fn show_catalog_now(
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Before anything else can see this catalog, and exactly once per open —
|
||||||
|
// the early return above is what makes it once. The job queue is durable,
|
||||||
|
// so a run that was killed mid-job left its row marked `Running` with
|
||||||
|
// nobody holding it; recovery hands those back to be resumed rather than
|
||||||
|
// lost, and drops jobs naming photographs that have since been deleted.
|
||||||
|
//
|
||||||
|
// Here rather than wherever a runner starts, because there is no owner
|
||||||
|
// column in `jobs`: a second recovery pass while a worker held a claim
|
||||||
|
// would take that claim away from it.
|
||||||
|
match dr_catalog::runner::recover(cat.connection()) {
|
||||||
|
Ok(r) if r.did_anything() => log::info!(
|
||||||
|
"job queue: {} interrupted job(s) resumed, \
|
||||||
|
{} for deleted photographs dropped",
|
||||||
|
r.reclaimed,
|
||||||
|
r.reaped
|
||||||
|
),
|
||||||
|
Ok(_) => {}
|
||||||
|
// Not surfaced. The queue is rebuildable like everything else in the
|
||||||
|
// catalog, and a grid that refuses to open because a background queue
|
||||||
|
// could not be tidied is the worse failure by some way.
|
||||||
|
Err(e) => log::warn!("job queue could not be recovered: {e}"),
|
||||||
|
}
|
||||||
|
|
||||||
// The sidebar before the grid, because the grid's badges read collection
|
// The sidebar before the grid, because the grid's badges read collection
|
||||||
// membership — the same order the scan's completion uses.
|
// membership — the same order the scan's completion uses.
|
||||||
crate::collections_ui::refresh_tree(window, coll_ctl, &cat);
|
crate::collections_ui::refresh_tree(window, coll_ctl, &cat);
|
||||||
|
|||||||
Reference in New Issue
Block a user