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.
This commit is contained in:
2026-08-09 12:22:31 +02:00
parent fbadf9afc8
commit c8bb08e661
29 changed files with 7193 additions and 234 deletions
+41
View File
@@ -0,0 +1,41 @@
//! TRACES: NFR-ARCH-4 | NFR-R5
//! Catalog errors.
//!
//! Typed and attached to the affected subject rather than panicking — a
//! corrupt row or a failed job marks one image and lets the batch continue.
/// Something went wrong talking to the catalog.
#[derive(Debug, thiserror::Error)]
pub enum CatalogError {
#[error("sqlite: {0}")]
Sqlite(#[from] rusqlite::Error),
/// The catalog was written by a newer build.
///
/// Opening it read-write would corrupt state this build cannot represent,
/// so the app refuses and says so (NFR-R5).
#[error("catalog schema v{found} is newer than this build supports (v{supported})")]
SchemaTooNew { found: i64, supported: i64 },
/// A scan could not reach a root at all.
///
/// Distinct from "files are missing": this aborts the scan *before* the
/// deletion sweep, because every folder would look unreached and the sweep
/// would delete the whole library (FR-CAT-9).
#[error("root {0} is unreachable; scan aborted without pruning")]
RootUnreachable(u64),
/// A smart collection whose selector references itself, directly or via
/// another collection.
#[error("collection {0} would form a cycle")]
CollectionCycle(u64),
#[error("no such collection: {0}")]
NoSuchCollection(u64),
#[error("malformed stored selector: {0}")]
BadSelector(String),
#[error("io: {0}")]
Io(String),
}
+392
View File
@@ -0,0 +1,392 @@
//! 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());
}
}
+374
View File
@@ -0,0 +1,374 @@
//! 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
//! - [`query`] — selectors compiled to indexed SQL, windowed for the grid
//! - [`jobs`] — the durable background work queue
//! - [`merge`] / [`sync`] — cross-device collection merging
//!
//! # 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 error;
pub mod jobs;
pub mod merge;
pub mod query;
pub mod scan;
pub mod schema;
pub mod sync;
pub use error::CatalogError;
pub use jobs::{Job, JobKind, Priority};
pub use merge::MergeReport;
pub use query::{Query, Sort};
pub use scan::{DirAction, DirState, EntryAction, ScanOutcome};
/// 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.
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.
pub fn for_span(seconds: i64) -> Self {
const DAY: i64 = 86_400;
match seconds {
s if s > 5 * 365 * DAY => Granularity::Year,
s if s > 90 * DAY => Granularity::Month,
s if s > 2 * DAY => Granularity::Day,
_ => Granularity::Hour,
}
}
}
/// 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)?;
schema::migrate(&conn)?;
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)?;
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
WHERE {} AND captured_at IS NOT 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)
}
/// 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;
assert_eq!(Granularity::for_span(10 * 365 * DAY), Granularity::Year);
assert_eq!(Granularity::for_span(120 * DAY), Granularity::Month);
assert_eq!(Granularity::for_span(10 * DAY), Granularity::Day);
assert_eq!(Granularity::for_span(3600), 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");
}
}
+527
View File
@@ -0,0 +1,527 @@
//! TRACES: FR-CAT-7 | FR-NC-9
//! Merging a remote catalog's collections into the local one.
//!
//! # Why this is a merge and not a copy
//!
//! The catalog file syncs to Nextcloud, and a device that finds a newer remote
//! copy must not simply replace its own — whichever device synced second would
//! lose everything the first did not have. So the remote file is downloaded to
//! a side path, `ATTACH`ed, and merged table by table.
//!
//! Row-level merging needs identities that are stable across devices, and
//! `collections.id INTEGER PRIMARY KEY` is not: two devices independently
//! allocate id 1 for different collections. Hence `collections.uuid`, which is
//! what everything here keys on. The integer id stays local and is never
//! compared across catalogs.
//!
//! # Conflict rule
//!
//! Per collection, by `revision` — a monotonic counter bumped on every local
//! edit — with `modified` timestamp only as a tiebreak. Comparing revisions
//! rather than mtimes means a device with a skewed clock cannot silently win
//! (the failure mode FR-NC-9 avoids for sidecars, applied here).
//!
//! Membership merges as a **set union**, not last-writer-wins: two devices
//! each adding different images to the same collection keep both sets. That
//! is almost always what the user meant, and the exception — a removal racing
//! an addition — resolves in favour of the addition, which is recoverable by
//! removing it again. Silently losing an addition is not.
//!
//! # Deletion
//!
//! A deleted collection leaves a tombstone (`deleted = 1`), because a merge
//! against a device that still holds it would otherwise resurrect it. The
//! tombstone carries a revision like any other edit, so deletion competes on
//! the same footing as a rename.
use rusqlite::Connection;
use crate::error::CatalogError;
/// How a collection differed between the two catalogs.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MergeVerdict {
/// Present only remotely — insert it.
InsertedFromRemote,
/// Remote revision is higher — take its fields.
UpdatedFromRemote,
/// Local revision is at least as high — keep ours.
KeptLocal,
/// Remote says deleted, and wins on revision.
DeletedByRemote,
}
/// What a merge did, for logging and for telling the user.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct MergeReport {
pub inserted: usize,
pub updated: usize,
pub kept_local: usize,
pub deleted: usize,
pub members_added: usize,
}
impl MergeReport {
/// Whether the local catalog changed, and so needs re-uploading.
pub fn local_changed(&self) -> bool {
self.inserted > 0 || self.updated > 0 || self.deleted > 0 || self.members_added > 0
}
/// Whether the local catalog holds anything the remote did not, and so
/// must be uploaded even if nothing was taken from the remote.
pub fn should_upload(&self) -> bool {
self.kept_local > 0 || self.local_changed()
}
}
/// Decide one collection, given both sides' revisions.
///
/// Split out from the SQL so the rule is testable on its own — it is the part
/// that decides whether a user loses a collection.
pub fn verdict(
local: Option<(i64, i64)>, // (revision, modified)
remote: (i64, i64),
remote_deleted: bool,
) -> MergeVerdict {
let (r_rev, r_mod) = remote;
match local {
None if remote_deleted => {
// A tombstone for something we never had. Recording it still
// matters: without it, a third device could reintroduce the
// collection through us.
MergeVerdict::DeletedByRemote
}
None => MergeVerdict::InsertedFromRemote,
Some((l_rev, l_mod)) => {
// Revision first; timestamp only to break an exact tie. Equal
// revisions with equal timestamps keep local, so a merge that
// changes nothing is stable and repeatable.
let remote_wins = r_rev > l_rev || (r_rev == l_rev && r_mod > l_mod);
if !remote_wins {
MergeVerdict::KeptLocal
} else if remote_deleted {
MergeVerdict::DeletedByRemote
} else {
MergeVerdict::UpdatedFromRemote
}
}
}
}
/// Merge collections and membership from an attached catalog.
///
/// The remote catalog must already be attached under the schema name
/// `remote_cat`; [`crate::Catalog::merge_attached_collections`] handles that.
///
/// Runs in one transaction: a merge either lands whole or not at all.
pub fn merge_collections(conn: &Connection) -> Result<MergeReport, CatalogError> {
let tx = conn.unchecked_transaction()?;
let mut report = MergeReport::default();
// ---- collections ------------------------------------------------------
{
let mut stmt = tx.prepare(
"SELECT r.uuid, r.name, r.parent_id, r.kind, r.selector_json,
r.created, r.revision, r.modified, r.deleted,
l.revision, l.modified
FROM remote_cat.collections r
LEFT JOIN main.collections l ON l.uuid = r.uuid",
)?;
struct Incoming {
uuid: String,
name: String,
kind: i64,
selector_json: Option<String>,
created: i64,
revision: i64,
modified: i64,
// No `deleted` field: the verdict already encodes it, and keeping
// both invites the two disagreeing.
verdict: MergeVerdict,
}
let rows: Vec<Incoming> = stmt
.query_map([], |r| {
let deleted: i64 = r.get(8)?;
let local_rev: Option<i64> = r.get(9)?;
let local_mod: Option<i64> = r.get(10)?;
let revision: i64 = r.get(6)?;
let modified: i64 = r.get(7)?;
Ok(Incoming {
uuid: r.get(0)?,
name: r.get(1)?,
kind: r.get(3)?,
selector_json: r.get(4)?,
created: r.get(5)?,
revision,
modified,
verdict: verdict(local_rev.zip(local_mod), (revision, modified), deleted != 0),
})
})?
.collect::<Result<_, _>>()?;
for row in rows {
match row.verdict {
MergeVerdict::KeptLocal => {
report.kept_local += 1;
}
MergeVerdict::InsertedFromRemote => {
tx.execute(
"INSERT INTO main.collections
(uuid, name, parent_id, kind, selector_json,
created, revision, modified, deleted)
VALUES (?1, ?2, NULL, ?3, ?4, ?5, ?6, ?7, 0)",
rusqlite::params![
row.uuid,
row.name,
row.kind,
row.selector_json,
row.created,
row.revision,
row.modified,
],
)?;
report.inserted += 1;
}
MergeVerdict::UpdatedFromRemote => {
tx.execute(
"UPDATE main.collections
SET name = ?2, kind = ?3, selector_json = ?4,
revision = ?5, modified = ?6, deleted = 0
WHERE uuid = ?1",
rusqlite::params![
row.uuid,
row.name,
row.kind,
row.selector_json,
row.revision,
row.modified,
],
)?;
report.updated += 1;
}
MergeVerdict::DeletedByRemote => {
// Tombstone rather than DELETE: the row must outlive the
// deletion or a third device reintroduces it.
tx.execute(
"INSERT INTO main.collections
(uuid, name, kind, created, revision, modified, deleted)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, 1)
ON CONFLICT(uuid) DO UPDATE SET
deleted = 1, revision = ?5, modified = ?6",
rusqlite::params![
row.uuid,
row.name,
row.kind,
row.created,
row.revision,
row.modified,
],
)?;
tx.execute(
"DELETE FROM main.collection_members
WHERE collection_id = (SELECT id FROM main.collections WHERE uuid = ?1)",
[&row.uuid],
)?;
report.deleted += 1;
}
}
}
}
// ---- membership -------------------------------------------------------
//
// Set union, keyed on (collection uuid, image content hash). The hash
// rather than the image id, for the same reason collections use a uuid:
// image ids are local. An image the remote has and we do not is skipped —
// it will join when a scan or sync catalogues it, and the next merge picks
// it up.
//
// Tombstoned collections are excluded, or a merge would repopulate a
// collection it had just deleted.
let added = tx.execute(
"INSERT OR IGNORE INTO main.collection_members(collection_id, image_id, position, added)
SELECT lc.id, li.id, rm.position, rm.added
FROM remote_cat.collection_members rm
JOIN remote_cat.collections rc ON rc.id = rm.collection_id
JOIN main.collections lc ON lc.uuid = rc.uuid AND lc.deleted = 0
JOIN remote_cat.images ri ON ri.id = rm.image_id
JOIN main.images li ON li.content_hash = ri.content_hash
WHERE ri.content_hash IS NOT NULL",
[],
)?;
report.members_added = added;
tx.commit()?;
Ok(report)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::schema;
#[test]
fn a_collection_we_lack_is_taken_from_remote() {
assert_eq!(
verdict(None, (1, 100), false),
MergeVerdict::InsertedFromRemote
);
}
#[test]
fn higher_remote_revision_wins() {
assert_eq!(
verdict(Some((3, 100)), (4, 50), false),
MergeVerdict::UpdatedFromRemote
);
}
#[test]
fn a_skewed_clock_cannot_beat_a_higher_local_revision() {
// The remote's timestamp is far in the future, but it has seen fewer
// edits. Revision decides, so the skewed device does not silently
// overwrite real work.
assert_eq!(
verdict(Some((9, 100)), (2, 999_999), false),
MergeVerdict::KeptLocal
);
}
#[test]
fn equal_revisions_break_on_timestamp() {
assert_eq!(
verdict(Some((3, 100)), (3, 200), false),
MergeVerdict::UpdatedFromRemote
);
assert_eq!(
verdict(Some((3, 200)), (3, 100), false),
MergeVerdict::KeptLocal
);
}
#[test]
fn an_identical_collection_is_stable() {
// Merging twice must not oscillate or report spurious changes.
assert_eq!(
verdict(Some((3, 100)), (3, 100), false),
MergeVerdict::KeptLocal
);
}
#[test]
fn deletion_competes_on_revision_like_any_other_edit() {
// Remote deleted it at revision 5; we renamed it at revision 4. The
// deletion is newer, so it wins.
assert_eq!(
verdict(Some((4, 100)), (5, 100), true),
MergeVerdict::DeletedByRemote
);
// But a stale deletion does not undo a newer local edit.
assert_eq!(
verdict(Some((6, 100)), (5, 100), true),
MergeVerdict::KeptLocal
);
}
#[test]
fn a_tombstone_for_something_we_never_had_is_recorded() {
// Otherwise this device could reintroduce the collection to a third.
assert_eq!(verdict(None, (2, 100), true), MergeVerdict::DeletedByRemote);
}
// ---- integration over two real catalogs ------------------------------
fn two_catalogs() -> Connection {
let c = Connection::open_in_memory().unwrap();
schema::configure(&c).unwrap();
schema::migrate(&c).unwrap();
// A second in-memory database standing in for the downloaded remote.
c.execute_batch("ATTACH ':memory:' AS remote_cat").unwrap();
let remote_schema = super::super::schema::v1_for_attached("remote_cat");
c.execute_batch(&remote_schema).unwrap();
c
}
fn add_image(c: &Connection, db: &str, id: i64, hash: &str) {
c.execute(
&format!(
"INSERT INTO {db}.roots(id, kind, label) VALUES (1, 'local', 'r')
ON CONFLICT(id) DO NOTHING"
),
[],
)
.unwrap();
c.execute(
&format!(
"INSERT INTO {db}.images(id, root_id, source_ref, content_hash, added_at)
VALUES (?1, 1, ?2, ?3, 0)"
),
rusqlite::params![id, format!("img{id}.CR3"), hash],
)
.unwrap();
}
fn add_collection(c: &Connection, db: &str, id: i64, uuid: &str, name: &str, rev: i64) {
c.execute(
&format!(
"INSERT INTO {db}.collections(id, uuid, name, kind, created, revision, modified)
VALUES (?1, ?2, ?3, 0, 0, ?4, ?4)"
),
rusqlite::params![id, uuid, name, rev],
)
.unwrap();
}
#[test]
fn disjoint_collections_from_two_devices_both_survive() {
// The property the whole design exists for: neither device loses work.
let c = two_catalogs();
add_collection(&c, "main", 1, "uuid-local", "Iceland", 1);
add_collection(&c, "remote_cat", 1, "uuid-remote", "Portugal", 1);
let report = merge_collections(&c).unwrap();
assert_eq!(report.inserted, 1);
let names: Vec<String> = c
.prepare("SELECT name FROM main.collections ORDER BY name")
.unwrap()
.query_map([], |r| r.get(0))
.unwrap()
.collect::<Result<_, _>>()
.unwrap();
assert_eq!(names, vec!["Iceland", "Portugal"]);
}
#[test]
fn membership_unions_rather_than_replacing() {
// Two devices each added a different image to the same collection.
let c = two_catalogs();
add_collection(&c, "main", 1, "shared", "Trip", 1);
add_collection(&c, "remote_cat", 1, "shared", "Trip", 1);
add_image(&c, "main", 1, "hash-a");
add_image(&c, "main", 2, "hash-b");
add_image(&c, "remote_cat", 1, "hash-b");
c.execute(
"INSERT INTO main.collection_members(collection_id, image_id, added) VALUES (1, 1, 0)",
[],
)
.unwrap();
c.execute(
"INSERT INTO remote_cat.collection_members(collection_id, image_id, added)
VALUES (1, 1, 0)",
[],
)
.unwrap();
merge_collections(&c).unwrap();
let n: i64 = c
.query_row("SELECT count(*) FROM main.collection_members", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(n, 2, "both devices' additions survive");
}
#[test]
fn membership_maps_across_devices_by_content_hash() {
// The same photograph carries different integer ids on each device.
// Keying on the id would attach the wrong image.
let c = two_catalogs();
add_collection(&c, "main", 1, "shared", "Trip", 1);
add_collection(&c, "remote_cat", 1, "shared", "Trip", 1);
add_image(&c, "main", 77, "same-photo");
add_image(&c, "remote_cat", 3, "same-photo");
c.execute(
"INSERT INTO remote_cat.collection_members(collection_id, image_id, added)
VALUES (1, 3, 0)",
[],
)
.unwrap();
merge_collections(&c).unwrap();
let img: i64 = c
.query_row("SELECT image_id FROM main.collection_members", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(img, 77, "resolved to the local id for the same photo");
}
#[test]
fn an_image_we_do_not_have_yet_is_skipped_not_errored() {
let c = two_catalogs();
add_collection(&c, "main", 1, "shared", "Trip", 1);
add_collection(&c, "remote_cat", 1, "shared", "Trip", 1);
add_image(&c, "remote_cat", 1, "not-here-yet");
c.execute(
"INSERT INTO remote_cat.collection_members(collection_id, image_id, added)
VALUES (1, 1, 0)",
[],
)
.unwrap();
let report = merge_collections(&c).unwrap();
assert_eq!(report.members_added, 0);
// It joins on a later merge, once a scan has catalogued the file.
}
#[test]
fn a_remote_deletion_does_not_resurrect_via_membership() {
let c = two_catalogs();
add_collection(&c, "main", 1, "doomed", "Old", 1);
add_image(&c, "main", 1, "hash-a");
add_image(&c, "remote_cat", 1, "hash-a");
c.execute(
"INSERT INTO remote_cat.collections(id, uuid, name, kind, created, revision, modified, deleted)
VALUES (1, 'doomed', 'Old', 0, 0, 5, 5, 1)",
[],
)
.unwrap();
c.execute(
"INSERT INTO remote_cat.collection_members(collection_id, image_id, added)
VALUES (1, 1, 0)",
[],
)
.unwrap();
let report = merge_collections(&c).unwrap();
assert_eq!(report.deleted, 1);
let n: i64 = c
.query_row("SELECT count(*) FROM main.collection_members", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(n, 0, "membership must not repopulate a deleted collection");
}
#[test]
fn merging_twice_changes_nothing_the_second_time() {
let c = two_catalogs();
add_collection(&c, "remote_cat", 1, "uuid-r", "Portugal", 1);
let first = merge_collections(&c).unwrap();
assert!(first.local_changed());
let second = merge_collections(&c).unwrap();
assert!(!second.local_changed(), "merge must be idempotent");
}
#[test]
fn keeping_local_still_marks_the_catalog_for_upload() {
// We hold something the remote does not, so the remote is stale even
// though we took nothing from it.
let c = two_catalogs();
add_collection(&c, "main", 1, "shared", "Renamed here", 5);
add_collection(&c, "remote_cat", 1, "shared", "Old name", 2);
let report = merge_collections(&c).unwrap();
assert_eq!(report.kept_local, 1);
assert!(report.should_upload());
}
}
+511
View File
@@ -0,0 +1,511 @@
//! TRACES: FR-CAT-4 | FR-CAT-6
//! Compiling a [`Selector`] into indexed SQL, and windowing the result.
//!
//! The UI never assembles SQL — it hands over a [`Query`] and receives a
//! window. Two properties matter:
//!
//! 1. **Nothing user-supplied is interpolated into SQL text.** Every value
//! binds as a parameter; `LIKE` patterns have their wildcards escaped.
//! 2. **Predicates hit indices.** Filtering 50k images must stay interactive
//! (FR-CAT-6), which means no expression over a column that would defeat
//! its index.
use dr_types::{Availability, ColourLabel, DateSelector, FlagState, Selector};
use rusqlite::types::Value;
/// What to show, and in what order.
#[derive(Debug, Clone)]
pub struct Query {
pub filter: Selector,
pub sort: Sort,
pub descending: bool,
}
impl Default for Query {
fn default() -> Self {
Query {
filter: Selector::All,
sort: Sort::CapturedAt,
descending: true,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Sort {
CapturedAt,
Added,
FileName,
Rating,
/// Manual order within a collection. Falls back to capture time where the
/// query is not scoped to one collection, since position is meaningless
/// outside it.
CollectionPosition,
}
impl Sort {
/// The ORDER BY fragment. Fixed strings — never user input.
///
/// Capture time sorts NULLs last regardless of direction: an image whose
/// EXIF has not been read yet (metadata_state 1) should not lead the grid
/// simply because its timestamp is unknown.
fn sql(self, descending: bool) -> &'static str {
match (self, descending) {
(Sort::CapturedAt, false) => {
"ORDER BY images.captured_at IS NULL, images.captured_at ASC, images.id ASC"
}
(Sort::CapturedAt, true) => {
"ORDER BY images.captured_at IS NULL, images.captured_at DESC, images.id DESC"
}
(Sort::Added, false) => "ORDER BY images.added_at ASC, images.id ASC",
(Sort::Added, true) => "ORDER BY images.added_at DESC, images.id DESC",
(Sort::FileName, false) => "ORDER BY images.source_ref ASC, images.id ASC",
(Sort::FileName, true) => "ORDER BY images.source_ref DESC, images.id DESC",
(Sort::Rating, false) => "ORDER BY v.rating ASC, images.id ASC",
(Sort::Rating, true) => "ORDER BY v.rating DESC, images.id DESC",
(Sort::CollectionPosition, false) => {
"ORDER BY cm.position IS NULL, cm.position ASC, images.captured_at ASC"
}
(Sort::CollectionPosition, true) => {
"ORDER BY cm.position IS NULL, cm.position DESC, images.captured_at DESC"
}
}
}
/// Whether this sort needs the default-version join.
fn needs_version(self) -> bool {
matches!(self, Sort::Rating)
}
/// Whether this sort needs a collection-membership join.
fn needs_membership(self) -> bool {
matches!(self, Sort::CollectionPosition)
}
}
/// A compiled WHERE clause plus its bound parameters.
///
/// Kept separate from the statement so `count` and `window` can share one
/// compilation.
#[derive(Debug, Default)]
pub struct Compiled {
pub where_sql: String,
pub params: Vec<Value>,
/// True if the filter depends on capture time, and therefore on EXIF that
/// a freshly scanned library may not have read yet. The UI surfaces this
/// rather than silently under-reporting.
pub needs_capture_time: bool,
}
/// Compile a selector to SQL against the `images` table.
///
/// `now` is passed rather than read from the clock so a rolling window is
/// reproducible in tests and consistent across one query.
pub fn compile(filter: &Selector, now: i64) -> Compiled {
let mut params = Vec::new();
let sql = if filter.is_unfiltered() {
"1".to_string()
} else {
emit(filter, now, &mut params)
};
Compiled {
where_sql: sql,
params,
needs_capture_time: filter.needs_capture_time(),
}
}
fn emit(s: &Selector, now: i64, p: &mut Vec<Value>) -> String {
match s {
Selector::All => "1".into(),
Selector::Collection(id) => {
p.push(Value::Integer(id.0 as i64));
format!(
"EXISTS (SELECT 1 FROM collection_members m
WHERE m.image_id = images.id AND m.collection_id = ?{})",
p.len()
)
}
Selector::Folder {
root,
path,
recursive,
} => {
p.push(Value::Integer(root.0 as i64));
let root_ix = p.len();
if *recursive {
// Prefix match on the folder path. `like_prefix` escapes the
// pattern metacharacters, so a folder literally named "50%"
// matches itself and not everything.
p.push(Value::Text(like_prefix(path)));
format!(
"images.folder_id IN (
SELECT id FROM folders
WHERE root_id = ?{root_ix}
AND (path = ?{p} OR path LIKE ?{p} || '/%' ESCAPE '\\'))",
p = p.len()
)
} else {
p.push(Value::Text(path.clone()));
format!(
"images.folder_id IN (
SELECT id FROM folders WHERE root_id = ?{root_ix} AND path = ?{})",
p.len()
)
}
}
Selector::DateRange(d) => emit_date(d, now, p),
Selector::Rating { min } => {
p.push(Value::Integer(*min as i64));
format!("{} >= ?{}", default_version_scalar("rating"), p.len())
}
Selector::Label(l) => {
p.push(Value::Integer(label_code(*l)));
format!("{} = ?{}", default_version_scalar("label"), p.len())
}
Selector::Flag(f) => {
p.push(Value::Integer(flag_code(*f)));
format!("{} = ?{}", default_version_scalar("flag"), p.len())
}
Selector::Keyword(k) => {
p.push(Value::Text(k.clone()));
format!(
"EXISTS (SELECT 1 FROM keywords kw
JOIN versions kv ON kv.id = kw.version_id
WHERE kv.image_id = images.id AND kw.keyword = ?{})",
p.len()
)
}
Selector::Camera(c) => {
p.push(Value::Text(c.clone()));
format!("images.camera = ?{}", p.len())
}
Selector::Lens(l) => {
p.push(Value::Text(l.clone()));
format!("images.lens = ?{}", p.len())
}
Selector::IsoRange { min, max } => {
p.push(Value::Integer(*min as i64));
let lo = p.len();
p.push(Value::Integer(*max as i64));
format!("images.iso BETWEEN ?{lo} AND ?{}", p.len())
}
Selector::Availability(a) => {
p.push(Value::Integer(availability_code(*a)));
format!("images.availability = ?{}", p.len())
}
Selector::Text(t) => {
// Substring over filename and keywords. A LIKE scan is adequate at
// 50k; if free text over title and description becomes a real
// workflow, FTS5 is the answer and it is additive.
p.push(Value::Text(format!("%{}%", escape_like(t))));
let ix = p.len();
format!(
"(images.source_ref LIKE ?{ix} ESCAPE '\\'
OR EXISTS (SELECT 1 FROM keywords kw
JOIN versions kv ON kv.id = kw.version_id
WHERE kv.image_id = images.id
AND kw.keyword LIKE ?{ix} ESCAPE '\\'))"
)
}
// An empty conjunction is vacuously true; an empty disjunction matches
// nothing. Both arise from a UI that lets every term be cleared, and
// conflating them would show the whole library when the user meant the
// opposite.
Selector::All_(v) if v.is_empty() => "1".into(),
Selector::Any(v) if v.is_empty() => "0".into(),
Selector::All_(v) => join(v, " AND ", now, p),
Selector::Any(v) => join(v, " OR ", now, p),
Selector::Not(inner) => format!("NOT ({})", emit(inner, now, p)),
}
}
fn join(items: &[Selector], op: &str, now: i64, p: &mut Vec<Value>) -> String {
let parts: Vec<String> = items.iter().map(|s| emit(s, now, p)).collect();
format!("({})", parts.join(op))
}
fn emit_date(d: &DateSelector, now: i64, p: &mut Vec<Value>) -> String {
match d {
DateSelector::Between { from, to } => {
p.push(Value::Integer(*from));
let lo = p.len();
p.push(Value::Integer(*to));
// Half-open, so adjacent ranges neither overlap nor gap.
format!(
"(images.captured_at >= ?{lo} AND images.captured_at < ?{})",
p.len()
)
}
DateSelector::Rolling { days } => {
let from = now - (*days as i64) * 86_400;
p.push(Value::Integer(from));
format!("images.captured_at >= ?{}", p.len())
}
DateSelector::CollectionSpan(id) => {
p.push(Value::Integer(id.0 as i64));
let ix = p.len();
format!(
"images.captured_at BETWEEN
(SELECT min(i2.captured_at) FROM images i2
JOIN collection_members m2 ON m2.image_id = i2.id
WHERE m2.collection_id = ?{ix})
AND (SELECT max(i2.captured_at) FROM images i2
JOIN collection_members m2 ON m2.image_id = i2.id
WHERE m2.collection_id = ?{ix})"
)
}
}
}
/// Rating, label, and flag live on the *default* version, not the image.
///
/// A correlated subquery rather than a join, so these compose inside `OR` and
/// `NOT` without the join multiplying rows.
fn default_version_scalar(col: &str) -> String {
format!(
"(SELECT dv.{col} FROM versions dv
WHERE dv.image_id = images.id AND dv.is_default = 1 LIMIT 1)"
)
}
/// Escape LIKE metacharacters so a literal `%` or `_` in user text matches
/// itself. Paired with `ESCAPE '\'` in every LIKE that uses it.
fn escape_like(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for c in s.chars() {
if matches!(c, '%' | '_' | '\\') {
out.push('\\');
}
out.push(c);
}
out
}
fn like_prefix(path: &str) -> String {
escape_like(path.trim_end_matches('/'))
}
fn label_code(l: ColourLabel) -> i64 {
match l {
ColourLabel::Red => 1,
ColourLabel::Yellow => 2,
ColourLabel::Green => 3,
ColourLabel::Blue => 4,
ColourLabel::Purple => 5,
}
}
fn flag_code(f: FlagState) -> i64 {
match f {
FlagState::Unflagged => 0,
FlagState::Pick => 1,
FlagState::Reject => 2,
}
}
fn availability_code(a: Availability) -> i64 {
match a {
Availability::MetadataOnly => 0,
Availability::Preview => 1,
Availability::Original => 2,
Availability::Offline => 3,
}
}
/// Build the full SELECT for a window of results.
///
/// Joins are added only where the sort needs them, so an unsorted-by-rating
/// grid query touches one table.
pub fn window_sql(q: &Query, compiled: &Compiled) -> String {
let mut joins = String::new();
if q.sort.needs_version() {
joins.push_str(" LEFT JOIN versions v ON v.image_id = images.id AND v.is_default = 1");
}
if q.sort.needs_membership() {
// Only meaningful when the filter scopes to one collection; elsewhere
// position is NULL and the sort falls through to capture time.
joins.push_str(" LEFT JOIN collection_members cm ON cm.image_id = images.id");
}
format!(
"SELECT images.id, images.source_ref, images.availability, images.captured_at, \
images.captured_offset, images.metadata_state \
FROM images{joins} WHERE {} {} LIMIT ? OFFSET ?",
compiled.where_sql,
q.sort.sql(q.descending)
)
}
/// Build the COUNT for the same filter.
pub fn count_sql(compiled: &Compiled) -> String {
format!("SELECT count(*) FROM images WHERE {}", compiled.where_sql)
}
#[cfg(test)]
mod tests {
use super::*;
use dr_types::{CollectionId, RootId};
#[test]
fn unfiltered_compiles_to_a_constant() {
let c = compile(&Selector::All, 0);
assert_eq!(c.where_sql, "1");
assert!(c.params.is_empty());
}
#[test]
fn empty_conjunction_and_disjunction_differ() {
// The distinction that decides whether clearing a filter shows
// everything or nothing.
assert_eq!(compile(&Selector::All_(vec![]), 0).where_sql, "1");
assert_eq!(compile(&Selector::Any(vec![]), 0).where_sql, "0");
}
#[test]
fn values_bind_rather_than_interpolate() {
// The injection guard: a hostile keyword must appear in params, never
// in SQL text.
let evil = "'; DROP TABLE images; --";
let c = compile(&Selector::Keyword(evil.into()), 0);
assert!(!c.where_sql.contains("DROP"));
assert_eq!(c.params, vec![Value::Text(evil.into())]);
}
#[test]
fn like_metacharacters_are_escaped() {
// A search for "50%" must not match everything containing "50".
let c = compile(&Selector::Text("50%".into()), 0);
assert_eq!(c.params, vec![Value::Text("%50\\%%".into())]);
assert!(c.where_sql.contains("ESCAPE"));
}
#[test]
fn a_backslash_in_search_text_is_itself_escaped() {
let c = compile(&Selector::Text("a\\b".into()), 0);
assert_eq!(c.params, vec![Value::Text("%a\\\\b%".into())]);
}
#[test]
fn rolling_window_resolves_against_supplied_now() {
// Passed in rather than read from the clock, so the window is stable
// across one query and reproducible in a test.
let now = 1_000_000i64;
let c = compile(
&Selector::DateRange(DateSelector::Rolling { days: 90 }),
now,
);
assert_eq!(c.params, vec![Value::Integer(now - 90 * 86_400)]);
}
#[test]
fn between_is_half_open() {
let c = compile(
&Selector::DateRange(DateSelector::Between { from: 10, to: 20 }),
0,
);
// Half-open so adjacent day buckets neither overlap nor leave a gap.
assert!(c.where_sql.contains(">= ?1"));
assert!(c.where_sql.contains("< ?2"));
}
#[test]
fn nested_composition_numbers_parameters_in_order() {
let s = Selector::All_(vec![
Selector::Rating { min: 4 },
Selector::Any(vec![
Selector::Camera("X-T5".into()),
Selector::Not(Box::new(Selector::Lens("XF 35".into()))),
]),
]);
let c = compile(&s, 0);
assert_eq!(
c.params,
vec![
Value::Integer(4),
Value::Text("X-T5".into()),
Value::Text("XF 35".into()),
]
);
assert!(c.where_sql.contains("?1"));
assert!(c.where_sql.contains("?2"));
assert!(c.where_sql.contains("?3"));
}
#[test]
fn recursive_folder_matches_the_folder_itself_and_below() {
let c = compile(
&Selector::Folder {
root: RootId(1),
path: "2026/08".into(),
recursive: true,
},
0,
);
// Both branches: the folder's own images and those in subfolders.
assert!(c.where_sql.contains("path = ?2"));
assert!(c.where_sql.contains("|| '/%'"));
}
#[test]
fn collection_span_binds_its_id_once_and_reuses_it() {
let c = compile(
&Selector::DateRange(DateSelector::CollectionSpan(CollectionId(7))),
0,
);
assert_eq!(c.params, vec![Value::Integer(7)]);
}
#[test]
fn capture_time_dependency_is_reported() {
let c = compile(&Selector::DateRange(DateSelector::Rolling { days: 7 }), 0);
assert!(c.needs_capture_time);
let c = compile(&Selector::Rating { min: 5 }, 0);
assert!(!c.needs_capture_time);
}
#[test]
fn capture_sort_puts_unknown_timestamps_last_in_both_directions() {
// An image whose EXIF has not been read yet must not lead the grid
// just because its timestamp is NULL.
assert!(Sort::CapturedAt.sql(true).contains("IS NULL"));
assert!(Sort::CapturedAt.sql(false).contains("IS NULL"));
}
#[test]
fn window_sql_joins_only_when_the_sort_needs_it() {
let c = compile(&Selector::All, 0);
let plain = window_sql(
&Query {
filter: Selector::All,
sort: Sort::CapturedAt,
descending: true,
},
&c,
);
assert!(!plain.contains("JOIN"));
let rated = window_sql(
&Query {
filter: Selector::All,
sort: Sort::Rating,
descending: true,
},
&c,
);
assert!(rated.contains("JOIN versions"));
}
}
+254
View File
@@ -0,0 +1,254 @@
//! TRACES: FR-CAT-1 | FR-CAT-9 | NFR-P1
//! Incremental scan: the local analogue of ETag pruning.
//!
//! Nextcloud propagates ETags up the tree, so one request proves a whole
//! library unchanged (ARCH §8.4). A filesystem offers no such guarantee — a
//! directory's mtime moves when its *direct* entries change and not when a
//! grandchild does, so there is no cheap "did anything below here change"
//! probe.
//!
//! Local scan therefore prunes at each level rather than at the root: one
//! metadata probe per directory when nothing changed, instead of one per file.
//! A 50k-image library in ~2k folders costs 2k probes, which is the difference
//! between meeting and missing NFR-P1 on SAF.
//!
//! This module holds the decision logic and the deletion-sweep rules; walking
//! an actual directory belongs to the platform layer, which supplies
//! [`DirState`] and [`DirEntry`].
use dr_types::FormatFilter;
/// What a directory looked like when last scanned, and what it looks like now.
///
/// Both fields are cheap to obtain: one `stat` locally, one
/// `DocumentsContract` metadata query on SAF.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct DirState {
pub mtime: i64,
/// Direct children, files and directories alike.
///
/// mtime alone misses a delete-and-create inside one timestamp tick, and
/// coarse-granularity providers widen that window. The count does not
/// close the hole — a paired add and remove moves neither — but a bare add
/// or remove moves the count, and those are far commoner.
pub entry_count: u32,
}
/// One entry from a directory listing.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DirEntry {
pub name: String,
pub is_dir: bool,
pub size: u64,
pub mtime: i64,
}
/// What the scanner should do with a directory, before listing it.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DirAction {
/// Contents unchanged. Skip the listing, but still recurse into known
/// children — without upward propagation, a deep change is invisible from
/// here.
RecurseOnly,
/// List and reconcile, then recurse.
ListAndRecurse,
}
/// Decide whether a directory needs listing.
pub fn classify_dir(stored: Option<DirState>, current: DirState) -> DirAction {
match stored {
Some(s) if s == current => DirAction::RecurseOnly,
_ => DirAction::ListAndRecurse,
}
}
/// What reconciling one listed entry against the catalog implies.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EntryAction {
/// Not catalogued. Insert at `metadata_state = 1` and queue EXIF.
Insert,
/// Catalogued and unchanged. The common case, and it must cost nothing.
Unchanged,
/// Size or mtime moved: re-read metadata, rebuild the thumbnail, and drop
/// the content hash, which is no longer valid.
Changed,
/// Recognised but not a format the user asked to scan for.
Ignored,
}
/// What the catalog already holds for a source.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct KnownFile {
pub size: u64,
pub mtime: i64,
}
/// Classify one listed file.
pub fn classify_entry(
entry: &DirEntry,
known: Option<KnownFile>,
formats: &FormatFilter,
) -> EntryAction {
if !formats.allows_name(&entry.name) {
return EntryAction::Ignored;
}
match known {
None => EntryAction::Insert,
Some(k) if k.size == entry.size && k.mtime == entry.mtime => EntryAction::Unchanged,
Some(_) => EntryAction::Changed,
}
}
/// Outcome of a scan, which decides whether pruning may run.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ScanOutcome {
/// Every reachable folder was visited.
Complete,
/// The user cancelled. Partial state is valid — jobs are resumable — but
/// unvisited folders must not be read as deleted.
Cancelled,
/// The root itself could not be opened: drive unplugged, SAF grant
/// revoked, share unmounted.
RootUnreachable,
/// Some subtree failed while the root was fine.
PartialFailure,
}
impl ScanOutcome {
/// Whether the deletion sweep may run.
///
/// **The most dangerous decision in the catalog.** The sweep deletes every
/// folder not reached by this scan's generation. After an incomplete scan
/// that is most of the library, so it runs only on `Complete`.
///
/// FR-CAT-9 draws exactly this line: a source *proven absent* may leave
/// the catalog; a source merely *unreachable* is marked offline and kept,
/// with its ratings and edits intact.
pub fn may_prune(self) -> bool {
matches!(self, ScanOutcome::Complete)
}
}
#[cfg(test)]
mod tests {
use super::*;
use dr_types::Format;
const A: DirState = DirState {
mtime: 100,
entry_count: 5,
};
#[test]
fn unchanged_directory_is_not_listed() {
assert_eq!(classify_dir(Some(A), A), DirAction::RecurseOnly);
}
#[test]
fn a_never_seen_directory_is_listed() {
assert_eq!(classify_dir(None, A), DirAction::ListAndRecurse);
}
#[test]
fn changed_mtime_forces_a_listing() {
let now = DirState { mtime: 101, ..A };
assert_eq!(classify_dir(Some(A), now), DirAction::ListAndRecurse);
}
#[test]
fn entry_count_catches_what_mtime_misses() {
// A file added within the same timestamp tick: mtime is unchanged, so
// mtime alone would skip this directory and lose the new image.
let now = DirState {
mtime: 100,
entry_count: 6,
};
assert_eq!(classify_dir(Some(A), now), DirAction::ListAndRecurse);
}
#[test]
fn unchanged_file_costs_nothing() {
let e = DirEntry {
name: "IMG_0001.CR3".into(),
is_dir: false,
size: 30_000_000,
mtime: 500,
};
let known = KnownFile {
size: 30_000_000,
mtime: 500,
};
assert_eq!(
classify_entry(&e, Some(known), &FormatFilter::all()),
EntryAction::Unchanged
);
}
#[test]
fn a_resaved_file_is_reprocessed() {
let e = DirEntry {
name: "IMG_0001.CR3".into(),
is_dir: false,
size: 30_000_001,
mtime: 900,
};
let known = KnownFile {
size: 30_000_000,
mtime: 500,
};
assert_eq!(
classify_entry(&e, Some(known), &FormatFilter::all()),
EntryAction::Changed
);
}
#[test]
fn format_filter_excludes_unwanted_types() {
let jpeg = DirEntry {
name: "IMG_0001.JPG".into(),
is_dir: false,
size: 1,
mtime: 1,
};
assert_eq!(
classify_entry(&jpeg, None, &FormatFilter::raw_only()),
EntryAction::Ignored
);
assert_eq!(
classify_entry(&jpeg, None, &FormatFilter::all()),
EntryAction::Insert
);
}
#[test]
fn a_placeholder_is_catalogued_as_the_image_it_stands_for() {
// 121,785 of these in a real synced folder (ARCH §9.0). Each must
// enter the catalog as a CR2 marked offline, not be skipped as an
// unknown ".nextcloud" type.
let stub = DirEntry {
name: "_MG_4130.CR2.nextcloud".into(),
is_dir: false,
size: 1,
mtime: 1,
};
assert_eq!(
classify_entry(&stub, None, &FormatFilter::from_formats([Format::Cr2])),
EntryAction::Insert
);
}
#[test]
fn pruning_requires_a_complete_scan() {
assert!(ScanOutcome::Complete.may_prune());
}
#[test]
fn an_unreachable_root_never_prunes() {
// The guard that stops an unplugged drive from deleting the library:
// every folder would look unreached, so the sweep would take all of
// them (FR-CAT-9).
assert!(!ScanOutcome::RootUnreachable.may_prune());
assert!(!ScanOutcome::Cancelled.may_prune());
assert!(!ScanOutcome::PartialFailure.may_prune());
}
}
+353
View File
@@ -0,0 +1,353 @@
//! TRACES: FR-CAT-2 | NFR-R5
//! Schema definition and forward-only migrations.
//!
//! The catalog is an *index*, not a source of truth (ARCH §6.12) — it is
//! deletable and rebuildable from sources plus sidecars. That is what makes
//! migration failure survivable, and why the recovery path is the normal
//! mechanism rather than a last resort.
//!
//! Migrations are forward-only, transactional, and idempotent on retry
//! (NFR-R5). The app refuses to open a catalog newer than it understands
//! rather than corrupting it.
use rusqlite::Connection;
use crate::error::CatalogError;
/// Schema version this build writes and understands.
pub const SCHEMA_VERSION: i64 = 1;
/// Apply migrations up to [`SCHEMA_VERSION`].
///
/// Returns the version migrated from, so callers can log or back up before a
/// real migration (NFR-R2 requires a backup before schema change).
pub fn migrate(conn: &Connection) -> Result<i64, CatalogError> {
let from: i64 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?;
if from > SCHEMA_VERSION {
return Err(CatalogError::SchemaTooNew {
found: from,
supported: SCHEMA_VERSION,
});
}
if from == SCHEMA_VERSION {
return Ok(from);
}
// Each step runs in its own transaction so a failure leaves the catalog
// at a coherent version rather than half-migrated.
if from < 1 {
let tx = conn.unchecked_transaction()?;
tx.execute_batch(V1)?;
tx.pragma_update(None, "user_version", 1)?;
tx.commit()?;
}
Ok(from)
}
/// Connection setup applied on every open, migration or not.
///
/// WAL is required by NFR-R1: it survives power loss without corruption, and
/// it lets a background job write while the grid reads.
pub fn configure(conn: &Connection) -> Result<(), CatalogError> {
conn.pragma_update(None, "journal_mode", "WAL")?;
// NORMAL rather than FULL: with WAL this is durable across process death
// (which is what FR-PLAT-AND-3 cares about) and only risks the last
// transaction on power loss. The catalog is rebuildable; the sidecars are
// not, and they are written separately with their own fsync discipline.
conn.pragma_update(None, "synchronous", "NORMAL")?;
conn.pragma_update(None, "foreign_keys", true)?;
// A scan touching thousands of rows is transient; let SQLite spill to
// memory rather than materialising temp b-trees on disk.
conn.pragma_update(None, "temp_store", "MEMORY")?;
Ok(())
}
/// The v1 schema rewritten to target an attached database.
///
/// Needed because a downloaded remote catalog is `ATTACH`ed under its own
/// schema name before merging, and tests build one from scratch. SQLite has no
/// "create these tables over there" form, so the names are rewritten.
///
/// The rewrite is textual and therefore only as good as the naming discipline
/// in [`V1`]: every `CREATE TABLE`/`CREATE INDEX` must name its object
/// unqualified, which they do.
pub fn v1_for_attached(schema_name: &str) -> String {
V1.replace("CREATE TABLE ", &format!("CREATE TABLE {schema_name}."))
.replace("CREATE INDEX ", &format!("CREATE INDEX {schema_name}."))
.replace(
"CREATE UNIQUE INDEX ",
&format!("CREATE UNIQUE INDEX {schema_name}."),
)
// REFERENCES within an attached schema resolve to that schema already,
// so foreign keys need no rewriting — but the ON clause of an index
// does, and `CREATE INDEX x.name ON table` is the correct form.
}
const V1: &str = r#"
-- Roots -------------------------------------------------------------------
CREATE TABLE roots (
id INTEGER PRIMARY KEY,
kind TEXT NOT NULL, -- 'local' | 'saf' | 'remote'
grant_blob BLOB, -- SAF persisted permission; NULL on Linux
label TEXT NOT NULL,
last_seen INTEGER,
-- Bumped once per completed scan. Folders record the generation they were
-- reached in; anything older was not reached and no longer exists.
scan_generation INTEGER NOT NULL DEFAULT 0
);
-- Folders: the unit of change detection, local and remote alike -----------
CREATE TABLE folders (
id INTEGER PRIMARY KEY,
root_id INTEGER NOT NULL REFERENCES roots(id) ON DELETE CASCADE,
parent_id INTEGER REFERENCES folders(id) ON DELETE CASCADE,
path TEXT NOT NULL,
-- Remote: the propagating ETag that makes a no-op sync one request.
etag TEXT,
-- Local: directory mtime plus direct-entry count. mtime alone misses a
-- paired create+delete inside one timestamp tick; the count narrows that.
mtime INTEGER,
entry_count INTEGER,
scanned_generation INTEGER NOT NULL DEFAULT 0,
UNIQUE(root_id, path)
);
CREATE INDEX folders_parent ON folders(parent_id);
-- Images ------------------------------------------------------------------
CREATE TABLE images (
id INTEGER PRIMARY KEY,
root_id INTEGER NOT NULL REFERENCES roots(id) ON DELETE CASCADE,
folder_id INTEGER REFERENCES folders(id) ON DELETE CASCADE,
source_ref TEXT NOT NULL,
-- Expensive: requires reading the whole file. Computed only when
-- something needs it (import dedup, reconnect-by-hash), never in a scan.
content_hash TEXT,
format TEXT,
w INTEGER,
h INTEGER,
-- UTC seconds. NULL until EXIF is read, or if the file carries none.
captured_at INTEGER,
-- Minutes east of UTC. A photograph's timestamp is local to where it was
-- taken; storing UTC alone makes a Tokyo shoot span two days in Paris.
captured_offset INTEGER,
camera TEXT,
lens TEXT,
iso INTEGER,
aperture REAL,
shutter REAL,
availability INTEGER NOT NULL DEFAULT 0,
file_size INTEGER,
file_mtime INTEGER,
-- 0 = nothing, 1 = stat-only, 2 = full EXIF. The grid is usable at 1.
metadata_state INTEGER NOT NULL DEFAULT 0,
sidecar_mtime INTEGER,
added_at INTEGER NOT NULL,
UNIQUE(root_id, source_ref)
);
CREATE INDEX images_captured ON images(captured_at);
CREATE INDEX images_folder ON images(folder_id);
-- Partial: content_hash is NULL for most rows most of the time, and the
-- non-NULL subset is exactly what reconnect and dedup query.
CREATE INDEX images_hash ON images(content_hash) WHERE content_hash IS NOT NULL;
-- Versions ----------------------------------------------------------------
CREATE TABLE versions (
id INTEGER PRIMARY KEY,
image_id INTEGER NOT NULL REFERENCES images(id) ON DELETE CASCADE,
uuid TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
is_default INTEGER NOT NULL DEFAULT 0,
graph_hash TEXT,
rating INTEGER NOT NULL DEFAULT 0,
label INTEGER,
flag INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX versions_image ON versions(image_id);
CREATE TABLE keywords (
version_id INTEGER NOT NULL REFERENCES versions(id) ON DELETE CASCADE,
keyword TEXT NOT NULL,
PRIMARY KEY(version_id, keyword)
);
CREATE INDEX keywords_term ON keywords(keyword);
-- Remote mapping ----------------------------------------------------------
CREATE TABLE remote (
image_id INTEGER PRIMARY KEY REFERENCES images(id) ON DELETE CASCADE,
-- oc:fileid — stable across server-side rename and move, so a move is not
-- a re-download of 80 MB.
file_id INTEGER NOT NULL,
etag TEXT,
sync_state INTEGER NOT NULL DEFAULT 0,
remote_path TEXT
);
CREATE UNIQUE INDEX remote_file ON remote(file_id);
-- Collections -------------------------------------------------------------
CREATE TABLE collections (
id INTEGER PRIMARY KEY,
-- Device-independent identity. The integer id is local and collides
-- across devices; the UUID is what a cross-device merge keys on.
uuid TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
parent_id INTEGER REFERENCES collections(id) ON DELETE CASCADE,
kind INTEGER NOT NULL, -- 0 = manual, 1 = smart
selector_json TEXT, -- smart only
created INTEGER NOT NULL,
-- Monotonic per collection, bumped on every local edit. Merge compares
-- these rather than file mtimes, so a clock-skewed device cannot silently
-- win.
revision INTEGER NOT NULL DEFAULT 1,
modified INTEGER NOT NULL,
-- Tombstone. A deleted collection must outlive its deletion, or a merge
-- with a device that still has it would resurrect it.
deleted INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE collection_members (
collection_id INTEGER NOT NULL REFERENCES collections(id) ON DELETE CASCADE,
image_id INTEGER NOT NULL REFERENCES images(id) ON DELETE CASCADE,
position INTEGER, -- manual ordering; NULL = by capture time
added INTEGER NOT NULL,
PRIMARY KEY(collection_id, image_id)
);
CREATE INDEX members_image ON collection_members(image_id);
-- Cache -------------------------------------------------------------------
CREATE TABLE cache (
id INTEGER PRIMARY KEY,
version_id INTEGER REFERENCES versions(id) ON DELETE CASCADE,
image_id INTEGER REFERENCES images(id) ON DELETE CASCADE,
kind INTEGER NOT NULL, -- thumbnail | proxy | original
resolution INTEGER,
graph_hash TEXT,
path TEXT NOT NULL,
bytes INTEGER NOT NULL,
last_used INTEGER NOT NULL
);
CREATE INDEX cache_lru ON cache(last_used);
CREATE TABLE cache_rules (
id INTEGER PRIMARY KEY,
selector_json TEXT NOT NULL,
tier INTEGER NOT NULL,
priority INTEGER NOT NULL DEFAULT 0,
enabled INTEGER NOT NULL DEFAULT 1
);
CREATE TABLE image_cache (
image_id INTEGER PRIMARY KEY REFERENCES images(id) ON DELETE CASCADE,
tier_actual INTEGER NOT NULL DEFAULT 0,
-- Materialised rather than recomputed, so the grid can draw availability
-- badges without evaluating every rule for every visible cell.
tier_desired INTEGER NOT NULL DEFAULT 0,
bytes INTEGER NOT NULL DEFAULT 0,
last_used INTEGER,
pinned_by_rule INTEGER REFERENCES cache_rules(id) ON DELETE SET NULL
);
-- Jobs --------------------------------------------------------------------
CREATE TABLE jobs (
id INTEGER PRIMARY KEY,
kind INTEGER NOT NULL,
subject_id INTEGER,
priority INTEGER NOT NULL DEFAULT 0,
state INTEGER NOT NULL DEFAULT 0, -- 0=pending 1=running 2=failed
attempts INTEGER NOT NULL DEFAULT 0,
not_before INTEGER NOT NULL DEFAULT 0,
payload TEXT,
last_error TEXT,
-- Coalescing. Enqueueing the same work twice updates one row rather than
-- queueing it twice, which is what makes "enqueue on any change" safe to
-- call liberally.
UNIQUE(kind, subject_id)
);
CREATE INDEX jobs_ready ON jobs(state, priority DESC, not_before);
"#;
#[cfg(test)]
mod tests {
use super::*;
fn mem() -> Connection {
let c = Connection::open_in_memory().unwrap();
configure(&c).unwrap();
c
}
#[test]
fn migrate_creates_schema_at_current_version() {
let c = mem();
assert_eq!(migrate(&c).unwrap(), 0);
let v: i64 = c
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(v, SCHEMA_VERSION);
}
#[test]
fn migrate_is_idempotent() {
let c = mem();
migrate(&c).unwrap();
// Re-running must not error or duplicate anything — NFR-R5 requires
// idempotency on retry, since a migration can be interrupted.
assert_eq!(migrate(&c).unwrap(), SCHEMA_VERSION);
}
#[test]
fn refuses_a_catalog_from_a_newer_build() {
let c = mem();
migrate(&c).unwrap();
c.pragma_update(None, "user_version", SCHEMA_VERSION + 1)
.unwrap();
// Opening it read-write would corrupt data this build cannot
// represent. Refusing is the specified behaviour (NFR-R5).
assert!(matches!(
migrate(&c),
Err(CatalogError::SchemaTooNew { .. })
));
}
#[test]
fn foreign_keys_cascade_from_root_to_image() {
let c = mem();
migrate(&c).unwrap();
c.execute(
"INSERT INTO roots(id, kind, label) VALUES (1, 'local', 'test')",
[],
)
.unwrap();
c.execute(
"INSERT INTO images(id, root_id, source_ref, added_at) VALUES (1, 1, 'a.CR3', 0)",
[],
)
.unwrap();
c.execute("DELETE FROM roots WHERE id = 1", []).unwrap();
let n: i64 = c
.query_row("SELECT count(*) FROM images", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 0, "images must not outlive their root");
}
#[test]
fn job_uniqueness_coalesces_rather_than_duplicating() {
let c = mem();
migrate(&c).unwrap();
for _ in 0..5 {
c.execute(
"INSERT INTO jobs(kind, subject_id, priority) VALUES (1, 42, 0)
ON CONFLICT(kind, subject_id)
DO UPDATE SET priority = max(priority, excluded.priority)",
[],
)
.unwrap();
}
let n: i64 = c
.query_row("SELECT count(*) FROM jobs", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 1, "five enqueues of the same work is one job");
}
}
+252
View File
@@ -0,0 +1,252 @@
//! TRACES: FR-CAT-7 | FR-NC-9 | NFR-R1
//! Preparing the catalog file for upload, and taking in a remote one.
//!
//! # The hazard this module exists to handle
//!
//! A WAL-mode SQLite database is not one file. Committed transactions can live
//! in `catalog.sqlite-wal` with the main file lagging behind, so copying
//! `catalog.sqlite` alone uploads a **torn snapshot**: internally consistent as
//! of some older point, missing everything since. Worse, a naive copy taken
//! while a writer is mid-transaction can be structurally corrupt.
//!
//! So an upload never copies the live file. It runs a TRUNCATE checkpoint to
//! fold the WAL back into the main file, then uses SQLite's own backup API to
//! take a consistent snapshot — which serialises correctly against concurrent
//! writers rather than racing them.
//!
//! # What is actually synced
//!
//! Only collections merge (see [`crate::merge`]). The rest of the catalog is a
//! *local index* of *local* storage — folder mtimes, cache paths, job rows —
//! and copying another device's version of those in would be actively wrong.
//! The remote file is read for its collections and then discarded.
//!
//! This is why the catalog remains disposable in the ARCH §6.12 sense: nothing
//! here makes the local database authoritative for anything a rebuild could
//! not recover.
use std::path::{Path, PathBuf};
use rusqlite::Connection;
use crate::error::CatalogError;
use crate::merge::{self, MergeReport};
/// Schema name the downloaded remote catalog is attached under.
const REMOTE_SCHEMA: &str = "remote_cat";
/// Fold the WAL into the main database file.
///
/// TRUNCATE rather than PASSIVE: passive checkpointing gives up when a reader
/// holds the WAL open, which would leave recent commits out of the snapshot
/// without saying so.
pub fn checkpoint(conn: &Connection) -> Result<(), CatalogError> {
conn.pragma_update(None, "wal_checkpoint", "TRUNCATE")?;
Ok(())
}
/// Write a consistent snapshot of the catalog to `dest`, ready to upload.
///
/// Uses the backup API rather than a filesystem copy so the snapshot is
/// coherent even with writers active. Callers should still prefer a quiet
/// moment — this competes with background jobs for the write lock.
pub fn snapshot_for_upload(conn: &Connection, dest: &Path) -> Result<(), CatalogError> {
checkpoint(conn)?;
let mut out = Connection::open(dest)?;
let backup = rusqlite::backup::Backup::new(conn, &mut out)?;
// SQLite's own "copy everything" sentinel is -1, but rusqlite asserts a
// positive page count, so ask for more pages than a catalog will ever
// have. The effect is the same: one step, no interleaved writers, no
// progress callback. A 50k-image catalog is tens of megabytes.
backup.run_to_completion(i32::MAX, std::time::Duration::ZERO, None)?;
Ok(())
}
/// Whether a downloaded remote catalog is worth merging.
///
/// Cheap guard before attaching: a remote written by a newer build may contain
/// tables and columns this one cannot read, and attempting the merge would
/// fail mid-transaction rather than declining cleanly.
pub fn remote_is_mergeable(remote: &Path) -> Result<bool, CatalogError> {
let conn = Connection::open_with_flags(
remote,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
)?;
let v: i64 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?;
Ok(v <= crate::schema::SCHEMA_VERSION)
}
/// Attach a downloaded remote catalog, merge its collections, detach.
///
/// The remote file is opened **read-only** — this device never writes to
/// another device's catalog, it only reads collections out of it.
pub fn merge_remote(conn: &Connection, remote: &Path) -> Result<MergeReport, CatalogError> {
if !remote_is_mergeable(remote)? {
return Err(CatalogError::SchemaTooNew {
found: -1,
supported: crate::schema::SCHEMA_VERSION,
});
}
// Path binds as a parameter; ATTACH accepts one, so a path containing a
// quote cannot break out into SQL.
conn.execute(
&format!("ATTACH DATABASE ?1 AS {REMOTE_SCHEMA}"),
[remote.to_string_lossy().as_ref()],
)?;
let result = merge::merge_collections(conn);
// Detach even if the merge failed, or the next attempt errors with
// "database remote_cat is already in use".
let detach = conn.execute(&format!("DETACH DATABASE {REMOTE_SCHEMA}"), []);
if let Err(e) = detach {
log::warn!("failed to detach remote catalog: {e}");
}
result
}
/// Where the catalog snapshot and the downloaded remote live.
///
/// Both are transient working files, not the catalog itself, so they belong in
/// the cache directory rather than beside the live database.
#[derive(Debug, Clone)]
pub struct SyncPaths {
pub upload_snapshot: PathBuf,
pub downloaded_remote: PathBuf,
}
impl SyncPaths {
pub fn in_dir(cache_dir: &Path) -> Self {
SyncPaths {
upload_snapshot: cache_dir.join("catalog-upload.sqlite"),
downloaded_remote: cache_dir.join("catalog-remote.sqlite"),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::schema;
fn seeded(path: &Path) -> Connection {
let c = Connection::open(path).unwrap();
schema::configure(&c).unwrap();
schema::migrate(&c).unwrap();
c
}
#[test]
fn snapshot_captures_committed_data() {
let dir = tempdir();
let live = dir.join("catalog.sqlite");
let snap = dir.join("snap.sqlite");
let c = seeded(&live);
c.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES ('u1', 'Iceland', 0, 0, 1, 1)",
[],
)
.unwrap();
snapshot_for_upload(&c, &snap).unwrap();
// The snapshot must hold the row even though it was written after the
// database was created — the torn-file failure this guards against.
let s = Connection::open(&snap).unwrap();
let name: String = s
.query_row("SELECT name FROM collections", [], |r| r.get(0))
.unwrap();
assert_eq!(name, "Iceland");
}
#[test]
fn a_remote_from_a_newer_build_is_declined_not_attempted() {
let dir = tempdir();
let remote = dir.join("remote.sqlite");
let r = seeded(&remote);
r.pragma_update(None, "user_version", schema::SCHEMA_VERSION + 1)
.unwrap();
drop(r);
assert!(!remote_is_mergeable(&remote).unwrap());
let local = seeded(&dir.join("local.sqlite"));
assert!(matches!(
merge_remote(&local, &remote),
Err(CatalogError::SchemaTooNew { .. })
));
}
#[test]
fn merge_remote_round_trips_a_collection() {
let dir = tempdir();
let remote_path = dir.join("remote.sqlite");
{
let r = seeded(&remote_path);
r.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES ('u-remote', 'Portugal', 0, 0, 1, 1)",
[],
)
.unwrap();
checkpoint(&r).unwrap();
}
let local = seeded(&dir.join("local.sqlite"));
local
.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES ('u-local', 'Iceland', 0, 0, 1, 1)",
[],
)
.unwrap();
let report = merge_remote(&local, &remote_path).unwrap();
assert_eq!(report.inserted, 1);
let n: i64 = local
.query_row("SELECT count(*) FROM collections", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 2);
}
#[test]
fn the_remote_can_be_merged_twice_without_attach_conflict() {
// Detach must happen even on the failure path, or the second attempt
// errors with "database remote_cat is already in use".
let dir = tempdir();
let remote_path = dir.join("remote.sqlite");
{
let r = seeded(&remote_path);
r.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES ('u-remote', 'Portugal', 0, 0, 1, 1)",
[],
)
.unwrap();
checkpoint(&r).unwrap();
}
let local = seeded(&dir.join("local.sqlite"));
merge_remote(&local, &remote_path).unwrap();
let second = merge_remote(&local, &remote_path).unwrap();
assert!(!second.local_changed());
}
/// A scratch directory that cleans up with the test.
fn tempdir() -> PathBuf {
let base = std::env::temp_dir().join(format!(
"dr-catalog-test-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&base);
std::fs::create_dir_all(&base).unwrap();
base
}
}