Files
DarkRoom/core/dr-catalog/src/jobs.rs
T
dtourolle c8bb08e661 Add folder scan with format selection; validate A3 on a real library
Library setup as the user described it: pick a folder, choose which RAW
types to look for, scan recursively.

  dr-types::FormatFilter  the tick-box selection, seeing through VFS
                          placeholder suffixes so a dehydrated CR2 still
                          matches as a CR2
  dr-sync::scan           recursive walk, Depth:1 per directory, pruning
                          unchanged subtrees where the backend propagates
                          directory ETags

Verified against nextcloud.tourolle.paris (34.0.2) on a real library:

  browse root      32 entries, 98ms
  scan PhotosRaw   17,185 RAW files in 334 directories, 34.1s
                   (7,836 CR2 + 9,349 DNG)
  range read       262KB of a 21.5MB DNG in 119ms — 1.22% of the file,
                   and enough to read "Canon EOS 6D | ISO 100"

That last line is assumption A3 validated on real data. Cataloguing this
library by whole-file fetch would move roughly 370GB; the range path
moves a few MB.

Pruning is capability-gated rather than assumed: with per-entry ETags a
probe costs a request and proves nothing about children, so it is skipped
entirely. A test asserts zero probes in that case.

Still unresolved: /core/preview returns 400 for every parameter
combination tried, including on a JPEG the server reports as having a
preview. Not a request-shape bug — it fails identically bare. Recorded
rather than worked around; ARCH §6.7 already treats server previews as
opportunistic, so nothing depends on it.
2026-08-09 12:22:31 +02:00

393 lines
13 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.
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<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,
_ => 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<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.
/// 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> {
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<i64>>(2)?,
r.get::<_, i64>(3)?,
r.get::<_, i64>(4)?,
r.get::<_, Option<String>>(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<usize, CatalogError> {
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());
}
}