Files
DarkRoom/core/dr-catalog/src/lib.rs
T
dtourolleandClaude Opus 5 846a249156 Drain the queue that nothing has ever drained
`jobs` has been a complete durable work queue since the catalog was
written, and nothing has ever taken a job out of it. `claim_next`,
`complete`, `fail` and `recover_orphaned` had no callers outside their own
tests; `enqueue` had three. So the table grew one row per photograph and
kept it forever, and FR-PLAT-AND-3's resumability was a property of code
that never ran.

`runner` is the missing half. It owns no thread, no clock and no policy,
and that is the whole design: on Android the process does not decide when
background work may run. WorkManager does, subject to Doze, battery saver
and FR-NC-6's network constraints, and it revokes permission mid-job by
calling onStopped(). So the runner exposes `run_one` — claim, run, record —
and `drain`, which repeats it against a budget, a deadline and a
cancellation flag the host owns. A `Worker.doWork()` with ten minutes calls
drain with a deadline; a desktop idle pass calls it with none. That is the
seam the Android service plugs into, and it needs no Android to test.

Handlers are supplied from above, because the catalog knows what needs
doing and nothing about how: a thumbnail needs a decoder and a fetch needs
a network stack, neither of which belongs under core/dr-catalog. A runner
claims only kinds some handler declares, so a queue holding work this
device cannot do is left alone rather than failed five times.

Four outcomes, and only two of them are the job's fault. Done deletes the
row; Retry backs off; Abandon gives up now, for a failure no retry can fix;
Interrupted releases the claim with its attempt refunded and ends the
drain, because the host stopped rather than the job — five backgroundings
in a row must not mark good work as failed. Process death is the fifth and
cannot report itself, which is what `recover` is for.

Recovery is called from `show_catalog_now`, which is the one place a
catalog is opened for a session and already returns early if one is open.
It has to be exactly once and before any worker starts: there is no owner
column, so a second pass while a worker held a claim would take it away.
The attempt a dead claim consumed is deliberately kept — 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.

The tests cover claiming under contention twice over: sequentially across
two connections, and with four threads on four connections against one
catalog on disk, asserting every job ran exactly once. Plus completion,
backoff, giving up, abandoning, interruption, budget, deadline,
cancellation, and a job orphaned by a simulated crash being reclaimed and
run once rather than lost or repeated.

Not wired to a handler yet, and deliberately not: the only enqueue site
the app actually reaches is the remote scan's, whose thumbnails are already
served by the async grid worker, and `walk`'s two sites are reachable only
from the scan_local example. Inventing a handler to make the plumbing look
used is how a requirement comes to read as covered by code that does not
implement it.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-29 23:31:47 +02:00

564 lines
21 KiB
Rust

//! TRACES: FR-CAT-2 | FR-CAT-4 | FR-CAT-6 | NFR-P1
//! The catalog: a rebuildable index over the library.
//!
//! Not a source of truth. Sidecars next to the images hold the authoritative
//! edit state (ARCH §6.12), and this file is deletable at any time — rebuilt
//! by rescanning sources and reading sidecars. That inversion is deliberate:
//! darktable maintains both a database and sidecars while achieving the
//! reliability of neither.
//!
//! # What lives here
//!
//! - [`schema`] — tables and forward-only migrations
//! - [`scan`] — incremental discovery that prunes unchanged directories
//! - [`walk`] — those decisions driven against real storage, local or SAF
//! - [`query`] — selectors compiled to indexed SQL, windowed for the grid
//! - [`collections`] — the collection tree and membership the UI edits
//! - [`keywords`] — the keyword vocabulary and what it is assigned to
//! - [`faces`] — detected faces, the people they belong to, and who said so
//! - [`bursts`] — frames that are one moment, grouped so they judge as one
//! - [`jobs`] — the durable background work queue
//! - [`runner`] — the thing that drains it, driven by whoever owns the thread
//! - [`trash`] — soft delete to a folder, then permanent delete
//! - [`merge`] / [`sync`] — cross-device merging of collections and keywords
//!
//! # The one thing everything is designed around
//!
//! **Work is proportional to what changed, or to what the user is looking at —
//! never to library size.** A 50k-image library that has not changed costs one
//! metadata probe per folder to verify (§scan), no thumbnails to regenerate
//! (§jobs coalescing), and no rule evaluation per grid cell (materialised
//! `tier_desired`).
use std::path::Path;
use dr_types::{Availability, ImageId};
use rusqlite::Connection;
pub mod bursts;
pub mod cache;
pub mod collections;
pub mod dedup;
pub mod error;
pub mod face_shard;
pub mod faces;
pub mod jobs;
pub mod keywords;
pub mod merge;
pub mod query;
pub mod rating;
pub mod runner;
pub mod scan;
pub mod schema;
pub mod sync;
pub mod trash;
pub mod walk;
pub use cache::{Budget, Cache, DEFAULT_BUDGET_BYTES};
pub use collections::{Collection, CollectionKind, TreeRow};
pub use dedup::{seen_by_content, seen_by_metadata, set_content_hash};
pub use error::CatalogError;
pub use face_shard::{FaceShardStore, SharedFace};
pub use faces::{Calibration, DetectedFace, Face, FaceId, Person, PersonId};
pub use jobs::{Job, JobKind, Priority};
pub use keywords::{Coverage, Keyword, KeywordId, SelectionKeyword};
pub use merge::MergeReport;
pub use query::{Query, Sort};
pub use rating::{Judgement, MAX_RATING};
// Not `runner::Budget`: `cache::Budget` already owns that name here and
// means something else entirely (bytes on disk, not jobs in a slot).
// Callers spell the work budget `runner::Budget`, where it is unambiguous.
pub use runner::{DrainReport, JobHandler, Outcome, Runner};
pub use scan::{DirAction, DirState, EntryAction, ScanOutcome};
pub use trash::{TrashedImage, TRASH_DIR};
pub use walk::{ensure_root, mark_root_offline, scan_root, RootKind, ScanProgress, ScanReport};
/// One row of the library grid.
///
/// Exactly what a cell draws and nothing more — no join per cell, and
/// availability reads a materialised column rather than evaluating cache rules
/// (ARCH §9.5).
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GridRow {
pub id: ImageId,
pub name: String,
pub availability: Availability,
/// UTC seconds. `None` until EXIF has been read.
pub captured_at: Option<i64>,
/// Minutes east of UTC, for rendering the photographer's local time.
pub captured_offset: Option<i32>,
/// 0 = nothing, 1 = stat-only, 2 = full EXIF.
pub metadata_state: u8,
}
/// A count of images in one time bucket, for the timeline scrubber.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TimeBucket {
/// UTC seconds at the bucket's start.
pub start: i64,
pub count: u32,
}
/// Time bucket size.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Granularity {
Year,
Month,
Day,
Hour,
}
impl Granularity {
/// SQLite `strftime` format that collapses a timestamp to this bucket.
///
/// Applied to **local** time, not UTC: "everything from 3 August" means
/// the photographer's 3 August, which is why `captured_offset` is stored
/// alongside the UTC timestamp.
/// Public so a caller that must build its own bucketing query — one
/// joining collection membership, say — buckets identically to
/// [`Catalog::timeline_range`] rather than reimplementing the format.
pub fn strftime(self) -> &'static str {
match self {
Granularity::Year => "%Y",
Granularity::Month => "%Y-%m",
Granularity::Day => "%Y-%m-%d",
Granularity::Hour => "%Y-%m-%dT%H",
}
}
/// A sensible bucket size for a span of seconds, so the UI need not guess.
///
/// # Chosen by how many bars it produces, not by fixed cut-offs
///
/// This used to be four thresholds on the span, which reads sensibly and
/// behaves badly under zoom. Each zoom step halves the span, so the bar
/// count halves with it until a threshold is crossed — a fifteen-year
/// library went 15 bars, 8, then 46, 23, 11, and finally *6*. Zooming in
/// made the picture coarser, which is the opposite of what zooming is for.
///
/// So the choice is made on the axis's terms: of the four bucket sizes,
/// take the one whose bar count comes nearest [`Self::TARGET_BARS`]. The
/// count then stays in the same neighbourhood at every zoom level, and
/// each step in genuinely shows finer structure rather than the same
/// structure drawn wider.
///
/// Nearest in *ratio*, not in difference: the counts available for a given
/// span are orders of magnitude apart — a span is either about 4 years or
/// about 48 months — and on a linear measure the larger count always looks
/// further away, which would bias every choice towards too few bars.
pub fn for_span(seconds: i64) -> Self {
Self::for_bucket(seconds.max(1) / Self::TARGET_BARS)
}
/// The calendar unit nearest a bucket of `seconds`, for *labelling* one.
///
/// Split out from [`Self::for_span`] because the axis no longer buckets by
/// calendar unit at all — it divides the visible span into a fixed number
/// of equal bins (see `LibrarySettings::timeline_bars`). What is still
/// wanted is the unit a bin is closest to, so a bin of about a day is
/// labelled as a date and one of about a year as a year. Asked directly
/// rather than derived from the span, because the bin count is now the
/// user's rather than this module's target.
pub fn for_bucket(seconds: i64) -> Self {
let seconds = seconds.max(1) as f64;
// Finest first, so that when two options are equally far from the
// target the finer one wins: `min_by` keeps the first minimum it saw,
// and more detail is the better failure.
[
Granularity::Hour,
Granularity::Day,
Granularity::Month,
Granularity::Year,
]
.into_iter()
.min_by(|a, b| {
let cost = |g: Granularity| {
// How far off, measured multiplicatively: twice as long and
// half as long are equally wrong.
//
// Deliberately not clamped. A bucket shorter than the unit
// scores *worse* the coarser the unit, which is what makes an
// hour of photographs pick hourly bars instead of every option
// tying at "one bucket" and the coarsest winning.
(seconds / g.approx_seconds() as f64).ln().abs()
};
cost(*a)
.partial_cmp(&cost(*b))
// Ties cannot arise from real spans, but a NaN would; falling
// back to the coarser option keeps the axis drawable.
.unwrap_or(std::cmp::Ordering::Equal)
})
.unwrap_or(Granularity::Day)
}
/// How many bars the timeline wants across its axis.
///
/// Not a hard count — the bucket sizes are calendar units, so the actual
/// number lands where the calendar puts it. It is the figure the choice
/// aims at: enough bars that a busy fortnight is visibly busier than a
/// quiet one, few enough that each is wide enough to hit with a finger.
const TARGET_BARS: i64 = 40;
/// Nominal length of one bucket, for choosing between them.
///
/// Approximate on purpose: months and years vary and it does not matter
/// here, because this only ranks four options that are a factor of ~12 or
/// ~30 apart. The exact boundaries come from `strftime` on the real dates.
fn approx_seconds(self) -> i64 {
const DAY: i64 = 86_400;
match self {
Granularity::Year => 365 * DAY,
Granularity::Month => 30 * DAY,
Granularity::Day => DAY,
Granularity::Hour => 3600,
}
}
}
/// A connection to the catalog.
pub struct Catalog {
conn: Connection,
}
impl Catalog {
/// Open or create a catalog, migrating it forward if needed.
pub fn open(path: &Path) -> Result<Self, CatalogError> {
let conn = Connection::open(path)?;
schema::configure(&conn)?;
let from = schema::migrate(&conn)?;
// A migration adds a column; it cannot know what the value should be
// for rows that already existed. Backfilling on open is what stops
// those rows being silently partial.
for (what, n) in schema::backfill(&conn)? {
log::info!("backfilled {what} for {n} row(s) (schema was v{from})");
}
Ok(Catalog { conn })
}
/// An in-memory catalog, for tests and for a throwaway import preview.
pub fn in_memory() -> Result<Self, CatalogError> {
let conn = Connection::open_in_memory()?;
schema::configure(&conn)?;
schema::migrate(&conn)?;
schema::backfill(&conn)?;
Ok(Catalog { conn })
}
/// Escape hatch for modules that need raw access. Not part of the UI-facing
/// surface.
pub fn connection(&self) -> &Connection {
&self.conn
}
/// How many images match.
///
/// Returned alongside the first window so the grid can size its scrollbar
/// and paint in one round trip.
pub fn count(&self, q: &Query, now: i64) -> Result<usize, CatalogError> {
let c = query::compile(&q.filter, now);
let sql = query::count_sql(&c);
let n: i64 =
self.conn
.query_row(&sql, rusqlite::params_from_iter(c.params.iter()), |r| {
r.get(0)
})?;
Ok(n as usize)
}
/// Fetch one window of results.
///
/// Never returns the whole catalog: FR-CAT-4 requires memory bounded
/// independently of library size.
pub fn window(
&self,
q: &Query,
range: std::ops::Range<usize>,
now: i64,
) -> Result<Vec<GridRow>, CatalogError> {
let c = query::compile(&q.filter, now);
let sql = query::window_sql(q, &c);
let mut params = c.params.clone();
params.push(rusqlite::types::Value::Integer(range.len() as i64));
params.push(rusqlite::types::Value::Integer(range.start as i64));
let mut stmt = self.conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params_from_iter(params.iter()), |r| {
let source_ref: String = r.get(1)?;
let avail: i64 = r.get(2)?;
Ok(GridRow {
id: ImageId(r.get::<_, i64>(0)? as u64),
name: source_ref
.rsplit(['/', ':'])
.next()
.unwrap_or(&source_ref)
.to_string(),
availability: decode_availability(avail),
captured_at: r.get(3)?,
captured_offset: r.get::<_, Option<i64>>(4)?.map(|v| v as i32),
metadata_state: r.get::<_, i64>(5)? as u8,
})
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
/// Counts per time bucket, for the timeline scrubber.
///
/// One grouped aggregate over the `images_captured` index — not 50k rows
/// handed to the UI to bucket itself.
pub fn timeline(
&self,
q: &Query,
g: Granularity,
now: i64,
) -> Result<Vec<TimeBucket>, CatalogError> {
let c = query::compile(&q.filter, now);
// Bucketed in local time: captured_offset is minutes east of UTC, and
// NULL falls back to UTC rather than dropping the row.
let sql = format!(
"SELECT min(captured_at) AS start,
count(*) AS n
FROM images
-- A shadowed JPEG is the same frame as its RAW; counting both
-- would double every paired shot in the histogram.
WHERE {} AND captured_at IS NOT NULL AND shadowed_by IS NULL
GROUP BY strftime('{}', captured_at + coalesce(captured_offset, 0) * 60,
'unixepoch')
ORDER BY start ASC",
c.where_sql,
g.strftime()
);
let mut stmt = self.conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params_from_iter(c.params.iter()), |r| {
Ok(TimeBucket {
start: r.get(0)?,
count: r.get::<_, i64>(1)? as u32,
})
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
/// Counts per time bucket, bounded to a date range.
///
/// What a zoomed timeline needs: [`timeline`](Self::timeline) always spans
/// the whole library, so zooming in would return the same coarse buckets
/// with the ends cropped rather than finer detail over a narrower span.
pub fn timeline_range(
&self,
q: &Query,
g: Granularity,
from: i64,
to: i64,
now: i64,
) -> Result<Vec<TimeBucket>, CatalogError> {
let c = query::compile(&q.filter, now);
let sql = format!(
"SELECT min(captured_at) AS start,
count(*) AS n
FROM images
WHERE {} AND captured_at IS NOT NULL AND shadowed_by IS NULL
AND captured_at >= ?{} AND captured_at <= ?{}
GROUP BY strftime('{}', captured_at + coalesce(captured_offset, 0) * 60,
'unixepoch')
ORDER BY start ASC",
c.where_sql,
c.params.len() + 1,
c.params.len() + 2,
g.strftime()
);
let mut params = c.params.clone();
params.push(rusqlite::types::Value::Integer(from));
params.push(rusqlite::types::Value::Integer(to));
let mut stmt = self.conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params_from_iter(params.iter()), |r| {
Ok(TimeBucket {
start: r.get(0)?,
count: r.get::<_, i64>(1)? as u32,
})
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
/// Merge a downloaded remote catalog's collections into this one.
///
/// See [`sync`] for why only collections cross over.
pub fn merge_remote_catalog(&self, remote: &Path) -> Result<MergeReport, CatalogError> {
sync::merge_remote(&self.conn, remote)
}
/// Write a consistent snapshot ready to upload.
pub fn snapshot_for_upload(&self, dest: &Path) -> Result<(), CatalogError> {
sync::snapshot_for_upload(&self.conn, dest)
}
}
fn decode_availability(v: i64) -> Availability {
match v {
1 => Availability::Preview,
2 => Availability::Original,
3 => Availability::Offline,
_ => Availability::MetadataOnly,
}
}
#[cfg(test)]
mod tests {
use super::*;
use dr_types::Selector;
fn seeded() -> Catalog {
let cat = Catalog::in_memory().unwrap();
let c = cat.connection();
c.execute(
"INSERT INTO roots(id, kind, label) VALUES (1, 'local', 'lib')",
[],
)
.unwrap();
// Three images across two days, one with no EXIF read yet.
for (id, name, captured, state) in [
(1i64, "a.CR3", Some(1_000_000i64), 2i64),
(2, "b.CR3", Some(1_100_000), 2),
(3, "c.CR3", None, 1),
] {
c.execute(
"INSERT INTO images(id, root_id, source_ref, captured_at, metadata_state, added_at)
VALUES (?1, 1, ?2, ?3, ?4, 0)",
rusqlite::params![id, name, captured, state],
)
.unwrap();
}
cat
}
#[test]
fn count_and_window_agree() {
let cat = seeded();
let q = Query::default();
assert_eq!(cat.count(&q, 0).unwrap(), 3);
assert_eq!(cat.window(&q, 0..10, 0).unwrap().len(), 3);
}
#[test]
fn window_is_bounded_by_the_requested_range() {
// FR-CAT-4: memory independent of catalog size.
let cat = seeded();
let rows = cat.window(&Query::default(), 0..2, 0).unwrap();
assert_eq!(rows.len(), 2);
}
#[test]
fn paging_covers_every_row_exactly_once() {
let cat = seeded();
let q = Query::default();
let mut seen = Vec::new();
for start in (0..3).step_by(2) {
seen.extend(cat.window(&q, start..start + 2, 0).unwrap());
}
let mut ids: Vec<u64> = seen.iter().map(|r| r.id.0).collect();
ids.sort_unstable();
assert_eq!(ids, vec![1, 2, 3]);
}
#[test]
fn an_image_without_capture_time_sorts_last_not_first() {
// Otherwise a freshly scanned library leads with whatever has not been
// read yet, which looks like corruption to the user.
let cat = seeded();
let rows = cat.window(&Query::default(), 0..10, 0).unwrap();
assert_eq!(rows.last().unwrap().id, ImageId(3));
}
#[test]
fn metadata_state_reaches_the_grid() {
// The grid needs it to distinguish "no photos on this date" from
// "EXIF not read yet" (FR-NC-6c's honesty principle).
let cat = seeded();
let rows = cat.window(&Query::default(), 0..10, 0).unwrap();
let pending = rows.iter().find(|r| r.id == ImageId(3)).unwrap();
assert_eq!(pending.metadata_state, 1);
}
#[test]
fn a_filter_narrows_the_count() {
let cat = seeded();
let q = Query {
filter: Selector::Text("a.CR3".into()),
..Default::default()
};
assert_eq!(cat.count(&q, 0).unwrap(), 1);
}
#[test]
fn timeline_buckets_and_skips_unread_images() {
let cat = seeded();
let buckets = cat
.timeline(&Query::default(), Granularity::Day, 0)
.unwrap();
// Two images with timestamps, one day apart in UTC; the third has no
// capture time and cannot be placed on a timeline at all.
let total: u32 = buckets.iter().map(|b| b.count).sum();
assert_eq!(total, 2);
}
#[test]
fn timeline_granularity_follows_the_span() {
const DAY: i64 = 86_400;
// Chosen by how many bars it makes, not by fixed cut-offs — see
// `for_span`. Ten years of yearly bars is ten bars, which says almost
// nothing about a library; monthly is 122, which is a shape.
assert_eq!(Granularity::for_span(10 * 365 * DAY), Granularity::Month);
assert_eq!(Granularity::for_span(120 * DAY), Granularity::Day);
assert_eq!(Granularity::for_span(10 * DAY), Granularity::Day);
assert_eq!(Granularity::for_span(3600), Granularity::Hour);
// The property the target exists for: zooming in never coarsens the
// axis. Under the old thresholds a fifteen-year library went 15 bars,
// then 8, then 46, 23, 11 — finer spans drawn with wider bars.
let mut span = 15 * 365 * DAY;
let mut previous = Granularity::for_span(span).approx_seconds();
for _ in 0..10 {
span /= 2;
let bucket = Granularity::for_span(span).approx_seconds();
assert!(
bucket <= previous,
"halving the span to {span}s coarsened the bucket \
from {previous}s to {bucket}s"
);
previous = bucket;
}
// And a span shorter than any bucket still picks the finest, rather
// than every option tying at one bar and the coarsest winning.
assert_eq!(Granularity::for_span(60), Granularity::Hour);
assert_eq!(Granularity::for_span(1), Granularity::Hour);
}
#[test]
fn names_are_derived_for_both_paths_and_saf_ids() {
let cat = Catalog::in_memory().unwrap();
let c = cat.connection();
c.execute(
"INSERT INTO roots(id, kind, label) VALUES (1, 'saf', 'tree')",
[],
)
.unwrap();
c.execute(
"INSERT INTO images(id, root_id, source_ref, added_at)
VALUES (1, 1, 'primary:DCIM/Camera/IMG_1.CR3', 0)",
[],
)
.unwrap();
let rows = cat.window(&Query::default(), 0..10, 0).unwrap();
assert_eq!(rows[0].name, "IMG_1.CR3");
}
}