Wires dr-face to dr-catalog: a background sweep that reads the proxy the grid already built, detects, aligns, embeds and stores, then a clustering pass that turns those embeddings into suggested people. Detection runs on the Large thumbnail tier and nowhere else. That is what makes the feature affordable -- a browsed library has already paid for its proxies, so face indexing adds no RAW decode that was not already happening -- and it is why an image whose proxy is missing is skipped rather than fetched: requesting one here would put face indexing on the network path FR-CULL-8 keeps it off. The sweep keeps no cursor. It asks the catalog what is missing, so it resumes after process death with no repeated work beyond the in-flight image, and cancelling is dropping the receiver. recluster writes only the suggested half. Confirmed faces go in as anchors and come back untouched, and a cluster of one stays nameless -- naming every stray face would fill the People view with noise the user then has to dismiss. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
404 lines
14 KiB
Rust
404 lines
14 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,
|
|
/// 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 {
|
|
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)
|
|
}
|
|
}
|
|
|
|
/// 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());
|
|
}
|
|
}
|