The claim was a deferred transaction around a SELECT and an UPDATE, and under a single connection that is fine. Under two it is not what it looks like: the SELECT takes only a read lock, the UPDATE tries to upgrade, and in WAL a worker that read the same snapshot as another gets SQLITE_BUSY_SNAPSHOT on its write. That is not an error a busy handler can retry away — the fix is to roll back and start over — so the queue was "safe" only in the sense that the loser failed loudly instead of taking a job someone else was holding. `UPDATE jobs SET state = 1, attempts = attempts + 1 WHERE id = (SELECT ...) RETURNING ...` is one statement and so one implicit transaction that takes the write lock immediately. Two workers serialise, the loser waits out its busy timeout, and neither can see a row the other already holds. The existing tests are unchanged by it, because from one connection the two forms are indistinguishable — which is exactly why it was never noticed. The rest is the surface a runner has to have and did not: - `claim_next_matching` takes only kinds a worker can actually do. Without it a device with no connector claims `FetchOriginal`, fails it, and pays five wakeups and five backoffs per photograph to reach a conclusion known before it started. Filtering after a claim cannot work: the claim has already marked the row running. - `abandon` gives up now, for failures no retry can fix. `fail` uses it for its own MAX_ATTEMPTS branch, so there is one statement that ends a job. - `release` hands a claim back with its attempt refunded, for a worker that is being stopped rather than a job that is going wrong. `attempts` stands in for the owner column the table does not have: it is bumped by every claim, so a stale worker's release matches nothing and changes nothing. - `reap_orphan_subjects` deletes jobs whose photograph is gone. Coalescing keeps the table one row per unit of work and nothing ever shrank it when the work stopped existing. `ScanFolder` is excluded because its subject is a folder id, and joining that against `images` deletes by coincidence of numbering — hence `JobKind::subject_is_image`, and `JobKind::ALL` so the next kind added cannot quietly fall out of the filter. - `counts` is the number a foreground service's notification is built from. One behaviour change worth stating: a kind this build does not recognise is now parked with an error rather than read as `ExtractMetadata`. The old `unwrap_or` would have run a job of an unknown kind as some arbitrary known one, which is worse than not running it at all. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
872 lines
31 KiB
Rust
872 lines
31 KiB
Rust
//! TRACES: FR-CAT-3 | NFR-ARCH-2 | FR-PLAT-AND-3
|
|
//! The background work queue.
|
|
//!
|
|
//! Jobs live in the catalog, so they survive process death — routine on
|
|
//! Android rather than exceptional (FR-PLAT-AND-3). Two properties carry the
|
|
//! design:
|
|
//!
|
|
//! - **Coalescing.** `UNIQUE(kind, subject_id)` makes enqueueing idempotent,
|
|
//! so every code path that notices a change can just enqueue and let the
|
|
//! table absorb the redundancy.
|
|
//! - **Priority shared with the GPU scheduler** (ARCH §5.3), so one notion of
|
|
//! urgency governs the whole app and visible work always preempts bulk work.
|
|
//!
|
|
//! Nothing in this file runs a job. [`crate::runner`] is the other half — the
|
|
//! one that claims from this table, does the work through a handler, and
|
|
//! reports back. Worth knowing because for a long time it did not exist: every
|
|
//! producer called [`enqueue`] and nothing ever called [`claim_next`], so the
|
|
//! table only ever grew.
|
|
|
|
use rusqlite::{Connection, OptionalExtension};
|
|
|
|
use crate::error::CatalogError;
|
|
|
|
/// What a job does.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
#[repr(i64)]
|
|
pub enum JobKind {
|
|
/// Recursive incremental scan from a folder (§scan).
|
|
ScanFolder = 0,
|
|
/// Promote an image from stat-only to full EXIF.
|
|
ExtractMetadata = 1,
|
|
/// Build or rebuild a thumbnail.
|
|
Thumbnail = 2,
|
|
/// A sidecar on disk is newer than what the catalog read.
|
|
ReadSidecar = 3,
|
|
/// Flush a local edit to its sidecar. Debounced, never per slider tick.
|
|
WriteSidecar = 4,
|
|
/// Whole-file hash. On demand only — import dedup, reconnect-by-hash.
|
|
ContentHash = 5,
|
|
/// Range-extract an embedded preview from a remote file (FR-NC-3).
|
|
FetchPreview = 6,
|
|
/// Fetch a full original: pinned by rule, or explicitly asked for.
|
|
FetchOriginal = 7,
|
|
/// Detect and embed the faces in one image (FR-CULL-8).
|
|
///
|
|
/// One job does both, rather than splitting them: the proxy is already
|
|
/// decoded and in memory, and the natural unit of resumable work is one
|
|
/// photograph. Splitting would double the queue's row count for nothing.
|
|
///
|
|
/// Runs against the proxy tier, never a full decode — a library that has
|
|
/// been browsed has already paid for its proxies, so face indexing adds no
|
|
/// RAW decodes that were not already happening.
|
|
DetectFaces = 8,
|
|
}
|
|
|
|
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> {
|
|
Some(match v {
|
|
0 => JobKind::ScanFolder,
|
|
1 => JobKind::ExtractMetadata,
|
|
2 => JobKind::Thumbnail,
|
|
3 => JobKind::ReadSidecar,
|
|
4 => JobKind::WriteSidecar,
|
|
5 => JobKind::ContentHash,
|
|
6 => JobKind::FetchPreview,
|
|
7 => JobKind::FetchOriginal,
|
|
8 => JobKind::DetectFaces,
|
|
_ => return None,
|
|
})
|
|
}
|
|
|
|
/// Whether this job transfers over the network, and so is subject to the
|
|
/// metered-connection and charging constraints in FR-NC-6.
|
|
pub fn is_network(self) -> bool {
|
|
matches!(self, JobKind::FetchPreview | JobKind::FetchOriginal)
|
|
}
|
|
|
|
/// Whether `subject_id` names a row in `images`.
|
|
///
|
|
/// Every kind but one is per-photograph. `ScanFolder`'s subject is a
|
|
/// *folder*, and the two id spaces are unrelated — so anything that joins
|
|
/// `subject_id` against `images` has to exclude it, or it will read one
|
|
/// table's ids as another's and act on the answer.
|
|
pub fn subject_is_image(self) -> bool {
|
|
!matches!(self, JobKind::ScanFolder)
|
|
}
|
|
}
|
|
|
|
/// Scheduling class, matching the GPU tile scheduler (ARCH §5.3).
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
|
|
#[repr(i64)]
|
|
pub enum Priority {
|
|
/// Bulk work: metadata sweeps, rule-driven fetches, hashing.
|
|
Background = 0,
|
|
/// Just outside the viewport; the next image in culling.
|
|
Prefetch = 1,
|
|
/// Visible cells, and the image currently open.
|
|
///
|
|
/// Strictly preempts background work. Without this, scrolling during a
|
|
/// bulk thumbnail pass misses its frame budget — the common case, not an
|
|
/// edge case (NFR-ARCH-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.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
#[repr(i64)]
|
|
pub enum JobState {
|
|
Pending = 0,
|
|
Running = 1,
|
|
Failed = 2,
|
|
}
|
|
|
|
/// A job ready to run.
|
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
|
pub struct Job {
|
|
pub id: i64,
|
|
pub kind: JobKind,
|
|
pub subject_id: Option<i64>,
|
|
pub priority: Priority,
|
|
pub attempts: i64,
|
|
pub payload: Option<String>,
|
|
}
|
|
|
|
/// Give up after this many attempts and attach the error to the subject.
|
|
///
|
|
/// One corrupt file must not stall the queue behind endless retries
|
|
/// (FR-RAW-4).
|
|
pub const MAX_ATTEMPTS: i64 = 5;
|
|
|
|
/// Backoff before retrying a failed job, in seconds.
|
|
///
|
|
/// Exponential, capped — a server that is down for an hour should not be
|
|
/// retried every second, and a transient decode failure should not wait an
|
|
/// hour.
|
|
pub fn backoff_seconds(attempts: i64) -> i64 {
|
|
const CAP: i64 = 300;
|
|
match attempts {
|
|
a if a <= 0 => 0,
|
|
a if a >= 9 => CAP,
|
|
a => (1i64 << (a - 1)).min(CAP),
|
|
}
|
|
}
|
|
|
|
/// Enqueue work, coalescing with any identical pending job.
|
|
///
|
|
/// Re-requesting at a higher priority *promotes* the existing row rather than
|
|
/// duplicating it, which is what lets the grid shout "this one is visible now"
|
|
/// about a job already queued in the background.
|
|
pub fn enqueue(
|
|
conn: &Connection,
|
|
kind: JobKind,
|
|
subject_id: Option<i64>,
|
|
priority: Priority,
|
|
payload: Option<&str>,
|
|
) -> Result<(), CatalogError> {
|
|
conn.execute(
|
|
"INSERT INTO jobs(kind, subject_id, priority, state, payload)
|
|
VALUES (?1, ?2, ?3, 0, ?4)
|
|
ON CONFLICT(kind, subject_id) DO UPDATE SET
|
|
priority = max(jobs.priority, excluded.priority),
|
|
-- A job that failed and is being re-requested deserves a fresh
|
|
-- start: the file may well have changed since it failed.
|
|
state = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.state END,
|
|
attempts = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.attempts END,
|
|
not_before = CASE WHEN jobs.state = 2 THEN 0 ELSE jobs.not_before END",
|
|
rusqlite::params![kind as i64, subject_id, priority as i64, payload],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Claim the next runnable job, highest priority first.
|
|
///
|
|
/// `now` is passed rather than read from the clock so backoff is testable.
|
|
pub fn claim_next(conn: &Connection, now: i64) -> Result<Option<Job>, CatalogError> {
|
|
claim(conn, now, None)
|
|
}
|
|
|
|
/// Claim the next runnable job of one of `kinds`.
|
|
///
|
|
/// What lets a runner take only the work it can actually do. A device with no
|
|
/// connector must leave `FetchOriginal` rows alone rather than claim them and
|
|
/// fail them five times each with backoff; a runner gated onto an unmetered
|
|
/// network (FR-NC-6) passes the local kinds only and leaves the transfers
|
|
/// where they are. Neither is expressible by filtering *after* a claim,
|
|
/// because the claim has already marked the row `Running`.
|
|
///
|
|
/// An empty list claims nothing, which is the honest reading of "there is
|
|
/// nothing this worker can do".
|
|
pub fn claim_next_matching(
|
|
conn: &Connection,
|
|
now: i64,
|
|
kinds: &[JobKind],
|
|
) -> Result<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
|
|
WHERE state = 0 AND not_before <= ?{filter}
|
|
ORDER BY priority DESC, id ASC
|
|
LIMIT 1)
|
|
RETURNING id, kind, subject_id, priority, attempts, payload"
|
|
);
|
|
|
|
let claimed = conn
|
|
.query_row(&sql, rusqlite::params_from_iter(args.iter()), |r| {
|
|
Ok((
|
|
r.get::<_, i64>(0)?,
|
|
r.get::<_, i64>(1)?,
|
|
r.get::<_, Option<i64>>(2)?,
|
|
r.get::<_, i64>(3)?,
|
|
r.get::<_, i64>(4)?,
|
|
r.get::<_, Option<String>>(5)?,
|
|
))
|
|
})
|
|
.optional()?;
|
|
|
|
let Some((id, kind, subject_id, priority, attempts, payload)) = claimed else {
|
|
return Ok(None);
|
|
};
|
|
|
|
let Some(kind) = JobKind::from_i64(kind) else {
|
|
// A row written by a build that knows a kind this one does not — a
|
|
// downgrade, or a catalog synced from a newer device. Running it as
|
|
// some other kind would be worse than not running it, so it is parked
|
|
// where the next claim will not see it again.
|
|
//
|
|
// Answering `None` understates what is queued for one pass. The
|
|
// alternative is a loop that keeps claiming the same unreadable row.
|
|
abandon(conn, id, &format!("unknown job kind {kind}"))?;
|
|
return Ok(None);
|
|
};
|
|
|
|
Ok(Some(Job {
|
|
id,
|
|
kind,
|
|
subject_id,
|
|
priority: Priority::from_i64(priority),
|
|
attempts,
|
|
payload,
|
|
}))
|
|
}
|
|
|
|
/// Job finished successfully.
|
|
pub fn complete(conn: &Connection, id: i64) -> Result<(), CatalogError> {
|
|
conn.execute("DELETE FROM jobs WHERE id = ?1", [id])?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Job failed. Reschedules with backoff, or gives up past [`MAX_ATTEMPTS`].
|
|
pub fn fail(conn: &Connection, job: &Job, now: i64, err: &str) -> Result<(), CatalogError> {
|
|
if job.attempts >= MAX_ATTEMPTS {
|
|
abandon(conn, job.id, err)
|
|
} else {
|
|
conn.execute(
|
|
"UPDATE jobs SET state = 0, not_before = ?2, last_error = ?3 WHERE id = ?1",
|
|
rusqlite::params![job.id, now + backoff_seconds(job.attempts), err],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
/// Give up on a job now, with no further retries.
|
|
///
|
|
/// For failures a retry cannot fix — the subject is gone, the payload is
|
|
/// unreadable, the format is one this build does not know. Walking the whole
|
|
/// retry ladder to reach a conclusion the first attempt already reached costs
|
|
/// five wakeups and five backoffs per photograph, which on a phone is the
|
|
/// difference the user notices.
|
|
///
|
|
/// The row is kept rather than deleted, because "this file failed and here is
|
|
/// why" is something the user is entitled to see (NFR-ARCH-4).
|
|
pub fn abandon(conn: &Connection, id: i64, err: &str) -> Result<(), CatalogError> {
|
|
conn.execute(
|
|
"UPDATE jobs SET state = 2, last_error = ?2 WHERE id = ?1",
|
|
rusqlite::params![id, err],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Put a claimed job back exactly as it was found.
|
|
///
|
|
/// For a worker that is being stopped rather than a job that is going wrong:
|
|
/// the platform revoked the slot, the user left the screen. The attempt the
|
|
/// claim consumed is given back, because nothing was learned about the file —
|
|
/// without that, five backgroundings in a row would mark good work as failed.
|
|
///
|
|
/// Guarded on the claim still being *this* claim. There is no owner column, so
|
|
/// `attempts` stands in for one: it is bumped by every claim, so the row only
|
|
/// still reads `state = 1` with the caller's own attempt number while nobody
|
|
/// else has taken it since. A worker that comes back after the recovery pass
|
|
/// handed its job to someone else therefore changes nothing, rather than
|
|
/// releasing a job another worker is in the middle of.
|
|
pub fn release(conn: &Connection, job: &Job) -> Result<(), CatalogError> {
|
|
conn.execute(
|
|
"UPDATE jobs SET state = 0, attempts = max(0, attempts - 1)
|
|
WHERE id = ?1 AND state = 1 AND attempts = ?2",
|
|
rusqlite::params![job.id, job.attempts],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Recover jobs orphaned by process death.
|
|
///
|
|
/// A row left `Running` has no owner — the process that claimed it is gone.
|
|
/// Called at startup, before any worker begins (FR-PLAT-AND-3).
|
|
///
|
|
/// The attempt the dead claim consumed is deliberately *not* refunded. A job
|
|
/// that takes the process down with it is indistinguishable from one that
|
|
/// fails, and the attempt counter is the only evidence that survives a death —
|
|
/// without it a poison-pill job is reclaimed and re-run forever.
|
|
pub fn recover_orphaned(conn: &Connection) -> Result<usize, CatalogError> {
|
|
let n = conn.execute("UPDATE jobs SET state = 0 WHERE state = 1", [])?;
|
|
Ok(n)
|
|
}
|
|
|
|
/// Delete jobs whose photograph is gone.
|
|
///
|
|
/// Coalescing keeps the table one row per unit of work, but nothing shrinks it
|
|
/// when the work stops existing: a library that has been culled carries a
|
|
/// thumbnail job for every photograph deleted since the last time anything
|
|
/// looked. Each one would be claimed, run, and failed five times.
|
|
///
|
|
/// Only kinds whose subject really is an image ([`JobKind::subject_is_image`])
|
|
/// are considered — `ScanFolder`'s subject is a folder id, and joining it
|
|
/// against `images` would delete jobs by coincidence of numbering.
|
|
///
|
|
/// There is no foreign key to do this instead. `jobs.subject_id` deliberately
|
|
/// references nothing: it means different tables for different kinds, and a
|
|
/// constraint that is right for eight of nine kinds is not a constraint.
|
|
pub fn reap_orphan_subjects(conn: &Connection) -> Result<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)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::schema;
|
|
|
|
fn db() -> Connection {
|
|
let c = Connection::open_in_memory().unwrap();
|
|
schema::configure(&c).unwrap();
|
|
schema::migrate(&c).unwrap();
|
|
c
|
|
}
|
|
|
|
#[test]
|
|
fn repeated_enqueue_coalesces() {
|
|
let c = db();
|
|
for _ in 0..10 {
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
}
|
|
let n: i64 = c
|
|
.query_row("SELECT count(*) FROM jobs", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(n, 1);
|
|
}
|
|
|
|
#[test]
|
|
fn re_enqueueing_at_higher_priority_promotes() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
// The grid scrolls this image into view.
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Interactive, None).unwrap();
|
|
|
|
let p: i64 = c
|
|
.query_row("SELECT priority FROM jobs", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(p, Priority::Interactive as i64);
|
|
}
|
|
|
|
#[test]
|
|
fn priority_never_regresses() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Interactive, None).unwrap();
|
|
// A background sweep must not demote work the user is waiting on.
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
|
|
let p: i64 = c
|
|
.query_row("SELECT priority FROM jobs", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(p, Priority::Interactive as i64);
|
|
}
|
|
|
|
#[test]
|
|
fn claim_takes_highest_priority_first() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
enqueue(&c, JobKind::Thumbnail, Some(2), Priority::Interactive, None).unwrap();
|
|
enqueue(&c, JobKind::Thumbnail, Some(3), Priority::Prefetch, None).unwrap();
|
|
|
|
let first = claim_next(&c, 0).unwrap().unwrap();
|
|
assert_eq!(first.subject_id, Some(2));
|
|
let second = claim_next(&c, 0).unwrap().unwrap();
|
|
assert_eq!(second.subject_id, Some(3));
|
|
}
|
|
|
|
#[test]
|
|
fn a_claimed_job_is_not_claimed_twice() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
assert!(claim_next(&c, 0).unwrap().is_some());
|
|
assert!(claim_next(&c, 0).unwrap().is_none());
|
|
}
|
|
|
|
#[test]
|
|
fn failure_backs_off_then_becomes_claimable_again() {
|
|
let c = db();
|
|
enqueue(
|
|
&c,
|
|
JobKind::FetchPreview,
|
|
Some(1),
|
|
Priority::Background,
|
|
None,
|
|
)
|
|
.unwrap();
|
|
let job = claim_next(&c, 100).unwrap().unwrap();
|
|
fail(&c, &job, 100, "network down").unwrap();
|
|
|
|
// Still backing off.
|
|
assert!(claim_next(&c, 100).unwrap().is_none());
|
|
// Past the backoff.
|
|
assert!(claim_next(&c, 100 + backoff_seconds(job.attempts))
|
|
.unwrap()
|
|
.is_some());
|
|
}
|
|
|
|
#[test]
|
|
fn a_persistently_failing_job_stops_retrying() {
|
|
let c = db();
|
|
enqueue(
|
|
&c,
|
|
JobKind::ExtractMetadata,
|
|
Some(1),
|
|
Priority::Background,
|
|
None,
|
|
)
|
|
.unwrap();
|
|
|
|
let mut now = 0;
|
|
for _ in 0..MAX_ATTEMPTS {
|
|
let job = claim_next(&c, now).unwrap().expect("should be claimable");
|
|
fail(&c, &job, now, "corrupt file").unwrap();
|
|
now += backoff_seconds(job.attempts);
|
|
}
|
|
|
|
// One corrupt file must not stall the queue forever (FR-RAW-4).
|
|
assert!(claim_next(&c, now + 100_000).unwrap().is_none());
|
|
let state: i64 = c
|
|
.query_row("SELECT state FROM jobs", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(state, JobState::Failed as i64);
|
|
}
|
|
|
|
#[test]
|
|
fn re_requesting_a_failed_job_gives_it_a_fresh_start() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
let mut now = 0;
|
|
for _ in 0..MAX_ATTEMPTS {
|
|
let job = claim_next(&c, now).unwrap().unwrap();
|
|
fail(&c, &job, now, "boom").unwrap();
|
|
now += backoff_seconds(job.attempts);
|
|
}
|
|
// The file changed on disk, so the old failure says nothing about it.
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Interactive, None).unwrap();
|
|
let job = claim_next(&c, now).unwrap().expect("retryable again");
|
|
assert_eq!(job.attempts, 1);
|
|
}
|
|
|
|
#[test]
|
|
fn orphaned_jobs_return_to_pending_on_restart() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
claim_next(&c, 0).unwrap().unwrap();
|
|
// Process dies here. Android does this routinely.
|
|
assert_eq!(recover_orphaned(&c).unwrap(), 1);
|
|
assert!(claim_next(&c, 0).unwrap().is_some());
|
|
}
|
|
|
|
#[test]
|
|
fn backoff_grows_then_caps() {
|
|
assert_eq!(backoff_seconds(0), 0);
|
|
assert_eq!(backoff_seconds(1), 1);
|
|
assert_eq!(backoff_seconds(3), 4);
|
|
assert_eq!(backoff_seconds(100), 300);
|
|
}
|
|
|
|
/// An image row, so a job has a subject that exists.
|
|
fn image(c: &Connection, id: i64) {
|
|
c.execute(
|
|
"INSERT INTO roots(id, kind, label) VALUES (1, 'local', '/lib')
|
|
ON CONFLICT DO NOTHING",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
c.execute(
|
|
"INSERT INTO images(id, root_id, source_ref, added_at) VALUES (?1, 1, ?2, 0)",
|
|
rusqlite::params![id, format!("/lib/{id}.CR3")],
|
|
)
|
|
.unwrap();
|
|
}
|
|
|
|
#[test]
|
|
fn a_runner_claims_only_the_kinds_it_names() {
|
|
// The property a filtered claim exists for: work this worker cannot do
|
|
// is left untouched — not claimed, not attempted, not failed.
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
enqueue(
|
|
&c,
|
|
JobKind::FetchOriginal,
|
|
Some(1),
|
|
Priority::Interactive,
|
|
None,
|
|
)
|
|
.unwrap();
|
|
|
|
// `FetchOriginal` is the higher priority and would be claimed first by
|
|
// an unfiltered claim. It is not this worker's to take.
|
|
let job = claim_next_matching(&c, 0, &[JobKind::Thumbnail])
|
|
.unwrap()
|
|
.expect("the thumbnail is claimable");
|
|
assert_eq!(job.kind, JobKind::Thumbnail);
|
|
assert!(claim_next_matching(&c, 0, &[JobKind::Thumbnail])
|
|
.unwrap()
|
|
.is_none());
|
|
|
|
let (state, attempts): (i64, i64) = c
|
|
.query_row(
|
|
"SELECT state, attempts FROM jobs WHERE kind = ?1",
|
|
[JobKind::FetchOriginal as i64],
|
|
|r| Ok((r.get(0)?, r.get(1)?)),
|
|
)
|
|
.unwrap();
|
|
assert_eq!(state, JobState::Pending as i64);
|
|
assert_eq!(attempts, 0);
|
|
}
|
|
|
|
#[test]
|
|
fn a_worker_that_can_do_nothing_claims_nothing() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
assert!(claim_next_matching(&c, 0, &[]).unwrap().is_none());
|
|
}
|
|
|
|
#[test]
|
|
fn releasing_a_claim_gives_the_attempt_back() {
|
|
// A stopped worker has learned nothing about the file, so the claim it
|
|
// is handing back must cost nothing.
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
|
|
let job = claim_next(&c, 0).unwrap().unwrap();
|
|
assert_eq!(job.attempts, 1);
|
|
release(&c, &job).unwrap();
|
|
|
|
let again = claim_next(&c, 0).unwrap().expect("claimable again at once");
|
|
assert_eq!(again.attempts, 1, "the release refunded the first attempt");
|
|
}
|
|
|
|
#[test]
|
|
fn releasing_a_job_someone_else_has_reclaimed_does_nothing() {
|
|
// `attempts` standing in for an owner column. A worker that comes back
|
|
// after the recovery pass handed its job to someone else must not
|
|
// release a claim that is no longer its to release — which would drop
|
|
// the live worker's job back into the queue to be run twice.
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
|
|
let stale = claim_next(&c, 0).unwrap().unwrap();
|
|
recover_orphaned(&c).unwrap();
|
|
let live = claim_next(&c, 0).unwrap().unwrap();
|
|
assert_eq!(live.attempts, 2);
|
|
|
|
release(&c, &stale).unwrap();
|
|
|
|
let (state, attempts): (i64, i64) = c
|
|
.query_row("SELECT state, attempts FROM jobs", [], |r| {
|
|
Ok((r.get(0)?, r.get(1)?))
|
|
})
|
|
.unwrap();
|
|
assert_eq!(
|
|
state,
|
|
JobState::Running as i64,
|
|
"the live claim still holds"
|
|
);
|
|
assert_eq!(attempts, 2, "and its attempt was not refunded for it");
|
|
}
|
|
|
|
#[test]
|
|
fn abandoning_skips_the_whole_retry_ladder() {
|
|
let c = db();
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
let job = claim_next(&c, 0).unwrap().unwrap();
|
|
|
|
abandon(&c, job.id, "not an image this build can read").unwrap();
|
|
|
|
assert!(claim_next(&c, 1_000_000).unwrap().is_none());
|
|
let (state, attempts, err): (i64, i64, String) = c
|
|
.query_row("SELECT state, attempts, last_error FROM jobs", [], |r| {
|
|
Ok((r.get(0)?, r.get(1)?, r.get(2)?))
|
|
})
|
|
.unwrap();
|
|
assert_eq!(state, JobState::Failed as i64);
|
|
assert_eq!(attempts, 1, "one attempt, not MAX_ATTEMPTS");
|
|
// Kept, not deleted: the user is entitled to see what failed and why.
|
|
assert!(err.contains("this build can read"));
|
|
}
|
|
|
|
#[test]
|
|
fn jobs_for_a_deleted_photograph_are_reaped() {
|
|
// Coalescing keeps the table one row per unit of work; nothing shrank
|
|
// it when the work stopped existing.
|
|
let c = db();
|
|
image(&c, 1);
|
|
image(&c, 2);
|
|
enqueue(&c, JobKind::Thumbnail, Some(1), Priority::Background, None).unwrap();
|
|
enqueue(&c, JobKind::Thumbnail, Some(2), Priority::Background, None).unwrap();
|
|
enqueue(
|
|
&c,
|
|
JobKind::ExtractMetadata,
|
|
Some(2),
|
|
Priority::Background,
|
|
None,
|
|
)
|
|
.unwrap();
|
|
|
|
c.execute("DELETE FROM images WHERE id = 2", []).unwrap();
|
|
assert_eq!(reap_orphan_subjects(&c).unwrap(), 2);
|
|
|
|
let left: i64 = c
|
|
.query_row("SELECT subject_id FROM jobs", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(left, 1);
|
|
}
|
|
|
|
#[test]
|
|
fn a_folder_scan_is_not_reaped_by_image_ids() {
|
|
// `ScanFolder`'s subject is a folder. Joining it against `images`
|
|
// would delete it whenever the numbering happened not to collide —
|
|
// which, on a fresh library, is almost always.
|
|
let c = db();
|
|
enqueue(
|
|
&c,
|
|
JobKind::ScanFolder,
|
|
Some(1),
|
|
Priority::Background,
|
|
Some("/lib/2024"),
|
|
)
|
|
.unwrap();
|
|
|
|
assert_eq!(reap_orphan_subjects(&c).unwrap(), 0);
|
|
assert!(claim_next(&c, 0).unwrap().is_some());
|
|
}
|
|
|
|
#[test]
|
|
fn a_job_kind_from_a_newer_build_is_parked_rather_than_guessed_at() {
|
|
// A catalog synced from a device running a later build. Running an
|
|
// unknown kind as some arbitrary known one is worse than not running
|
|
// it, and the old code silently read every unknown kind as
|
|
// `ExtractMetadata`.
|
|
let c = db();
|
|
c.execute(
|
|
"INSERT INTO jobs(kind, subject_id, priority, state) VALUES (99, 1, 0, 0)",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
|
|
assert!(claim_next(&c, 0).unwrap().is_none());
|
|
|
|
let (state, err): (i64, String) = c
|
|
.query_row("SELECT state, last_error FROM jobs", [], |r| {
|
|
Ok((r.get(0)?, r.get(1)?))
|
|
})
|
|
.unwrap();
|
|
assert_eq!(state, JobState::Failed as i64);
|
|
assert!(err.contains("99"), "{err}");
|
|
}
|
|
|
|
#[test]
|
|
fn counts_say_what_is_left() {
|
|
// What a foreground service's notification is built from: a number,
|
|
// without reading a hundred thousand rows to find it.
|
|
let c = db();
|
|
assert_eq!(counts(&c).unwrap(), Counts::default());
|
|
|
|
for id in 1..=3 {
|
|
enqueue(&c, JobKind::Thumbnail, Some(id), Priority::Background, None).unwrap();
|
|
}
|
|
let job = claim_next(&c, 0).unwrap().unwrap();
|
|
abandon(&c, job.id, "nope").unwrap();
|
|
claim_next(&c, 0).unwrap().unwrap();
|
|
|
|
let n = counts(&c).unwrap();
|
|
assert_eq!(n.pending, 1);
|
|
assert_eq!(n.running, 1);
|
|
assert_eq!(n.failed, 1);
|
|
assert_eq!(n.outstanding(), 2, "failed work is not outstanding work");
|
|
}
|
|
|
|
#[test]
|
|
fn every_kind_is_in_all() {
|
|
// `ALL` is what the reap builds its kind filter from, so a kind added
|
|
// to the enum and forgotten here would quietly stop being reaped.
|
|
for (i, kind) in JobKind::ALL.iter().enumerate() {
|
|
assert_eq!(
|
|
JobKind::from_i64(i as i64),
|
|
Some(*kind),
|
|
"ALL is out of step with the discriminants at {i}"
|
|
);
|
|
}
|
|
assert!(JobKind::from_i64(JobKind::ALL.len() as i64).is_none());
|
|
}
|
|
|
|
#[test]
|
|
fn only_a_folder_scan_has_a_non_image_subject() {
|
|
assert!(!JobKind::ScanFolder.subject_is_image());
|
|
for kind in JobKind::ALL.iter().filter(|k| **k != JobKind::ScanFolder) {
|
|
assert!(kind.subject_is_image(), "{kind:?}");
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn network_jobs_are_identifiable_for_metered_gating() {
|
|
// FR-NC-6: transfers respect unmetered-network and charging
|
|
// constraints; local work must not be gated by them.
|
|
assert!(JobKind::FetchOriginal.is_network());
|
|
assert!(JobKind::FetchPreview.is_network());
|
|
assert!(!JobKind::Thumbnail.is_network());
|
|
assert!(!JobKind::ExtractMetadata.is_network());
|
|
}
|
|
}
|