`persist` runs after every scan, for every photograph the scan listed. On a settled library that is the folders whose ETag changed -- a sidecar written there by a rating is enough -- so one relisted folder of 1,600 images is an ordinary pass, and a first scan is all 24,000. Per photograph it prepared four statements from their SQL (a folder lookup, the image upsert, the id read-back, the remote upsert) and then, after the commit, found the image again by path and enqueued its thumbnail job as an autocommitting statement of its own -- a commit per photograph, for rows that were almost all already queued. Now the statements are prepared once per pass, a folder's id is looked up once per folder rather than once per photograph in it, and the job is enqueued inside the transaction with the id already in hand. That also makes the job atomic with the row it points at, which is what the old ordering after the commit was trying to guarantee. `jobs::enqueue` uses a cached statement for the same reason. persist_bench on a copy of the reference catalog, CPU, best of runs: largest folder (1,589 images) 102-118 ms -> 10-13 ms whole library (23,582 images) 1.55-2.19 s -> 188-192 ms The fingerprint of images, remote, jobs and folders after the run is the same for both builds.
878 lines
32 KiB
Rust
878 lines
32 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> {
|
|
// Cached: a scan enqueues one per photograph it lists.
|
|
conn.prepare_cached(
|
|
"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",
|
|
)?
|
|
.execute(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());
|
|
}
|
|
}
|