//! 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. use rusqlite::Connection; 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, } impl JobKind { fn from_i64(v: i64) -> Option { 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, _ => 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) } } /// 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, } /// 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, pub priority: Priority, pub attempts: i64, pub payload: Option, } /// 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, 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. /// 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, CatalogError> { let tx = conn.unchecked_transaction()?; let job = tx .query_row( "SELECT id, kind, subject_id, priority, attempts, payload FROM jobs WHERE state = 0 AND not_before <= ?1 ORDER BY priority DESC, id ASC LIMIT 1", [now], |r| { Ok(( r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, Option>(2)?, r.get::<_, i64>(3)?, r.get::<_, i64>(4)?, r.get::<_, Option>(5)?, )) }, ) .ok(); let Some((id, kind, subject_id, priority, attempts, payload)) = job else { return Ok(None); }; tx.execute( "UPDATE jobs SET state = 1, attempts = attempts + 1 WHERE id = ?1", [id], )?; tx.commit()?; Ok(Some(Job { id, kind: JobKind::from_i64(kind).unwrap_or(JobKind::ExtractMetadata), subject_id, priority: match priority { 2 => Priority::Interactive, 1 => Priority::Prefetch, _ => Priority::Background, }, attempts: attempts + 1, 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 { conn.execute( "UPDATE jobs SET state = 2, last_error = ?2 WHERE id = ?1", rusqlite::params![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(()) } /// 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). pub fn recover_orphaned(conn: &Connection) -> Result { let n = conn.execute("UPDATE jobs SET state = 0 WHERE state = 1", [])?; Ok(n) } #[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); } #[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()); } }