Files
DarkRoom/ui/dr-ui/src/library/scan.rs
T
dtourolle 5da28584a4 Stop the scan queueing a thumbnail job per photograph
The reference catalog held 23,582 Thumbnail jobs, one per image, and
every scan re-coalesced all of them. Nothing has ever claimed that kind:
no JobHandler is registered for it on desktop or Android, and
dr_catalog::sync never merges another device's jobs in.

Thumbnails are owed by the store, not the queue. The grid's worker and
the thumbnail sweep both find their work by asking ThumbStore what it
lacks, and the store is shared between devices, so it is the only record
that knows another device already made one. A queue row was a second,
staler copy of that debt that grew with the library and was read by
nothing.

persist still writes the images and their remote identities in the one
transaction; it just no longer adds a row to jobs for each of them. The
two tests that asserted the rows existed become one that asserts a
repeated scan queues nothing.

Refs #73
2026-09-26 13:09:44 -04:00

1278 lines
49 KiB
Rust

//! Walking a remote library into the catalog, and pulling other devices'
//! judgements out of the sidecars the walk finds along the way.
use dr_catalog::Catalog;
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
use dr_types::FormatFilter;
use std::path::PathBuf;
use std::sync::mpsc::{Receiver, Sender};
use super::cells::total_images;
use super::sidecar::sidecar_path;
/// Progress and results from the scan worker.
#[derive(Debug)]
pub enum ScanMessage {
/// Directories walked so far, and images found.
Progress {
directories: usize,
pruned: usize,
images: usize,
},
/// The scan finished and the catalog is populated.
///
/// `found` counts what this scan *listed*, which on an incremental rescan
/// is only what changed — pruned directories contribute nothing. `total`
/// is what the catalog actually holds, which is what the grid shows.
/// Conflating them made a successful no-op rescan report "0 images" and
/// blank the library.
Done {
found: usize,
total: usize,
pruned: usize,
elapsed_ms: u64,
/// TRACES: FR-CAT-8 | FR-NC-9
/// Judgements this scan took *in* from other devices' sidecars.
///
/// Counted and reported rather than left to the log because it is the
/// only visible sign that a cull made elsewhere has arrived. A grid
/// that silently gains three hundred stars is indistinguishable from
/// one that has gone wrong.
judgements: usize,
},
/// The scan could not finish.
///
/// `offline` distinguishes "the server could not be reached" from "the
/// server refused", and it is carried here rather than re-derived because
/// the classification is only possible on the worker side: crossing the
/// channel flattens a [`dr_sync::RemoteError`] into a message, and no
/// amount of string matching on the far side can reliably recover it.
/// Without the flag a dead connection and a bad password produce the same
/// banner, which sends the user to re-enter a credential that was fine.
///
/// `lost_root` is the same idea one step further out, and it is carried
/// separately from `offline` rather than folded into it because the two
/// end differently. An offline library comes back when the network does,
/// with nothing asked of anyone; a library whose root cannot be opened
/// comes back only when someone restores access to it — a share put back
/// on the server, a drive plugged in, and in time a document tree granted
/// again once one can be (FR-PLAT-AND-2). Both show the same grid of what
/// is stored locally, and they must not offer the same explanation.
Failed {
message: String,
offline: bool,
lost_root: bool,
},
}
/// Run a scan on a worker thread, writing results into the catalog.
///
/// Returns the receiver the UI drains. The worker owns its own tokio runtime
/// and backend; nothing here touches the Slint event loop.
pub fn spawn_scan(
conn: Connection,
root: String,
filter: FormatFilter,
catalog_path: PathBuf,
) -> Receiver<ScanMessage> {
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let started = std::time::Instant::now();
if let Err(e) = run_scan(&tx, conn, root, filter, catalog_path, started) {
let _ = tx.send(ScanMessage::Failed {
message: e.message,
offline: e.offline,
lost_root: e.lost_root,
});
}
});
rx
}
/// A scan failure that still knows whether it was a connectivity failure.
///
/// The scan crosses a thread boundary, so the typed error cannot travel with
/// it; this carries the one bit that must survive.
pub(super) struct ScanFailure {
message: String,
offline: bool,
lost_root: bool,
}
impl ScanFailure {
/// A failure that is nothing to do with reachability — local I/O, a
/// runtime that would not start, a catalog that would not open.
fn local(message: impl std::fmt::Display) -> Self {
Self {
message: message.to_string(),
offline: false,
lost_root: false,
}
}
}
impl From<dr_sync::RemoteError> for ScanFailure {
fn from(e: dr_sync::RemoteError) -> Self {
Self {
offline: e.indicates_offline(),
lost_root: e.indicates_lost_root(),
message: e.to_string(),
}
}
}
/// TRACES: FR-PLAT-AND-2 | FR-CAT-9
/// Record that a library can no longer be opened, without losing it.
///
/// Called on the worker, before the failure crosses the channel, because this
/// is where the catalog handle is — and because the marking must be durable
/// whether or not anyone is left to draw a banner. A process killed between
/// the failure and the next launch must still come back knowing what it could
/// not reach.
///
/// Nothing is deleted. Every rating, every edit and every row stays exactly
/// where it was; what changes is that the images now say they are offline, so
/// the grid can show them as held-not-here rather than as ordinary
/// photographs whose thumbnails happen to be failing one at a time.
///
/// A root with no row yet is the first scan of a library that has never
/// succeeded, and there is nothing to mark — the failure alone is the whole
/// story, and the launch screen is where it is told.
pub(super) fn mark_library_offline(catalog: &Catalog, root: &str) {
let conn = catalog.connection();
let root_id: Option<i64> = conn
.query_row(
"SELECT id FROM roots WHERE label = ?1 AND kind = 'remote'",
[root],
|r| r.get(0),
)
.ok();
let Some(root_id) = root_id else {
log::info!("library {root} has no catalog root yet; nothing to mark offline");
return;
};
match dr_catalog::mark_root_offline(conn, dr_types::RootId(root_id as u64)) {
Ok(()) => log::warn!("library {root} is unreachable; its images are marked offline"),
Err(e) => log::error!("could not mark {root} offline: {e}"),
}
}
pub(super) fn run_scan(
tx: &Sender<ScanMessage>,
conn: Connection,
root: String,
filter: FormatFilter,
catalog_path: PathBuf,
started: std::time::Instant,
) -> Result<(), ScanFailure> {
if let Some(dir) = catalog_path.parent() {
std::fs::create_dir_all(dir)
.map_err(|e| ScanFailure::local(format!("creating {}: {e}", dir.display())))?;
}
let catalog = Catalog::open(&catalog_path).map_err(ScanFailure::local)?;
let rt = crate::net_runtime::build().map_err(ScanFailure::local)?;
rt.block_on(async {
// TRACES: FR-PLAT-AND-2 | FR-CAT-9
// Classified rather than flattened to a local failure, because the
// removed-card case never gets as far as a request: the folder
// connector checks its root when it is constructed, so a library on an
// ejected card fails here and not in the walk. Reported as an ordinary
// error it left the grid showing a healthy library of images that
// could no longer be opened, one silent thumbnail failure at a time.
let backend = match crate::remote::connect(&conn) {
Ok(b) => b,
Err(e) => {
if e.indicates_lost_root() {
mark_library_offline(&catalog, &root);
}
return Err(e.into());
}
};
// Stored folder ETags, so an unchanged subtree is skipped whole. On a
// first run this is empty and the walk is complete; on every run after
// it is what keeps cost proportional to what changed (ARCH §8.4).
let known = load_folder_etags(&catalog, &root);
let scanned = dr_sync::scan(&*backend, &RemotePath::new(&root), &filter, &known, |p| {
let _ = tx.send(ScanMessage::Progress {
directories: p.directories_listed,
pruned: p.directories_pruned,
images: p.images_found,
});
})
.await;
// TRACES: FR-PLAT-AND-2 | FR-CAT-9
// Written before the failure is reported, not after: the banner is a
// consequence of the catalog state and not the other way round, and a
// process that dies between the two must come back knowing.
let result = match scanned {
Ok(r) => r,
Err(e) => {
if e.indicates_lost_root() {
mark_library_offline(&catalog, &root);
}
return Err(e.into());
}
};
persist(&catalog, &root, &result).map_err(ScanFailure::local)?;
// TRACES: FR-CAT-8 | FR-NC-9
// Take in what other devices have judged. After `persist`, because a
// judgement lands on an image's version row and the image has to be in
// the catalog first — a sidecar seen in the same listing as a
// photograph this scan has only just discovered is the ordinary case
// on a library another device imported.
//
// Deliberately not fatal. A pull that fails leaves the ratings this
// device already had exactly where they were and the sidecar's ETag
// unrecorded, so the next scan tries again; failing the whole scan over
// it would blank a grid that was working.
let judgements = match pull_sidecars(&*backend, &catalog, &root, &result.sidecars).await {
Ok(n) => n,
Err(e) => {
log::warn!("reading sidecars from the library: {e}");
0
}
};
// Report what the catalog holds, not what this pass listed. An
// incremental rescan lists only what changed, so its own count is
// near zero on a healthy library.
let total = total_images(&catalog).unwrap_or(result.images.len());
let _ = tx.send(ScanMessage::Done {
found: result.images.len(),
total,
pruned: result.progress.directories_pruned,
elapsed_ms: started.elapsed().as_millis() as u64,
judgements,
});
Ok(())
})
}
/// TRACES: FR-CAT-8 | FR-CAT-9 | FR-NC-9
/// Take other devices' judgements out of the library's sidecars.
///
/// # Why this had to exist
///
/// A rating is written to the catalog and to the photograph's sidecar, and the
/// sidecar is the authoritative one (ARCH §6.12). Nothing ever read one back.
/// The scan indexed files, `derived_sync` exchanged thumbnails, collections
/// and keywords, `dr_catalog::merge` reconciled everything in the catalog
/// *except* `versions.rating` and `versions.flag`, and the single sidecar
/// reader ran when one photograph was opened in develop and handed its result
/// to the develop graph. `JobKind::ReadSidecar` had been declared for this
/// since the queue was written and was never enqueued or handled.
///
/// So judgements travelled outward only. The grid draws `versions.rating`, and
/// a cull done on another device could not reach it by any path the app had.
///
/// # What it costs
///
/// Nothing on a library nobody has edited. The sidecars come from listings the
/// walk was making anyway, an unchanged directory is pruned before it is
/// listed at all, and a sidecar whose ETag matches what this device last read
/// is skipped without a request. What is left is one GET per sidecar that
/// genuinely changed — which is the number of photographs somebody edited.
///
/// # Why a judgement can only be added, never withdrawn
///
/// `Version::merge`'s rule, applied here: zero is *unjudged*, and a device that
/// has never rated a frame is indistinguishable from one that deliberately
/// cleared it. Taking a remote zero over a local star would let a device that
/// was never involved erase an afternoon's culling. So a remote zero is
/// ignored, and the cost is that clearing a rating does not propagate.
///
/// Returns how many images gained a judgement.
pub(super) async fn pull_sidecars(
backend: &dyn RemoteBackend,
catalog: &Catalog,
root: &str,
seen: &[dr_sync::RemoteEntry],
) -> Result<usize, String> {
if seen.is_empty() {
return Ok(0);
}
let conn = catalog.connection();
let root_id: i64 = conn
.query_row(
"SELECT id FROM roots WHERE label = ?1 AND kind = 'remote'",
[root],
|r| r.get(0),
)
.map_err(|e| e.to_string())?;
let known = load_sidecar_etags(catalog, root_id);
let mut applied = 0usize;
for entry in seen {
let path = entry.path.as_str();
if known.get(path).is_some_and(|e| *e == entry.validator) {
continue;
}
let bytes = match backend.get(&RemoteId::Path(entry.path.clone()), None).await {
Ok(b) => b,
// Gone between the listing and the fetch, or not on this device and
// not worth materialising a whole library for. Neither is an error,
// and neither records an ETag — so the next scan tries again.
Err(e) => {
log::debug!("reading sidecar {path}: {e}");
continue;
}
};
let text = String::from_utf8_lossy(&bytes);
// TRACES: FR-CAT-13
// A standard XMP beside the photograph — Lightroom's, darktable's,
// anybody's — takes the other branch: reconciled field by field with
// the catalog winning, and a disagreement recorded for the reload
// the requirement asks to be offered. See `xmp_sync`.
if crate::xmp_sync::is_xmp(path) {
match crate::xmp_sync::take_in(conn, root_id, path, &text, now_secs()) {
Ok(taken) => {
applied += taken.changed;
if !taken.conflicts.is_empty() {
log::info!(
"xmp sidecar {path} disagrees with the catalog on {:?}; \
a reload is offered in Settings",
taken.conflicts
);
}
// Recorded only once it reached a photograph: a sidecar
// that arrived before its image is read again next time.
if taken.described > 0 {
record_sidecar_read(conn, root_id, path, &entry.validator);
}
}
Err(e) => log::warn!("xmp sidecar at {path} is unreadable ({e})"),
}
continue;
}
let mut sidecar = match dr_pipeline::Sidecar::parse(&text) {
Ok(s) => s,
// Unreadable is not empty. Recording the ETag would mean never
// looking at it again, and a build that understands it may be
// along; leaving it unrecorded costs one GET per scan and keeps
// the door open.
Err(e) => {
log::warn!("sidecar at {path} is unreadable ({e})");
continue;
}
};
sidecar.fuse_default_versions(None);
let Some(version) = sidecar.default_version() else {
// A sidecar with no default version — someone else's virtual copy
// and nothing more. Nothing to take, but it *was* read, so its
// ETag is recorded and it is not fetched again.
record_sidecar_read(conn, root_id, path, &entry.validator);
continue;
};
match apply_judgement(
conn,
root_id,
path,
version.rating,
version.flag,
version.label,
) {
Ok(n) => {
applied += n;
record_sidecar_read(conn, root_id, path, &entry.validator);
}
Err(e) => log::debug!("applying {path}: {e}"),
}
}
if applied > 0 {
log::info!("{applied} judgement(s) arrived from other devices");
}
Ok(applied)
}
/// What this device has already read, so a pull fetches only what changed.
pub(super) fn load_sidecar_etags(
catalog: &Catalog,
root_id: i64,
) -> std::collections::HashMap<String, dr_sync::Validator> {
let mut out = std::collections::HashMap::new();
let Ok(mut stmt) = catalog
.connection()
.prepare("SELECT path, etag FROM sidecars WHERE root_id = ?1 AND etag IS NOT NULL")
else {
return out;
};
if let Ok(rows) = stmt.query_map([root_id], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
}) {
for (path, etag) in rows.flatten() {
out.insert(path, dr_sync::Validator::new(etag));
}
}
out
}
/// Remember that this sidecar has been taken in at this ETag.
///
/// Written only after the judgement has landed, for the same reason
/// `dr_sync::scan` records a directory's ETag only after listing it: recording
/// it first would let a failure look like work already done, and the edit would
/// never be read again.
pub(super) fn record_sidecar_read(
conn: &rusqlite::Connection,
root_id: i64,
path: &str,
etag: &dr_sync::Validator,
) {
let done = conn.execute(
"INSERT INTO sidecars(root_id, path, etag, read_at) VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(root_id, path) DO UPDATE SET
etag = excluded.etag, read_at = excluded.read_at",
rusqlite::params![root_id, path, etag.as_str(), now_secs()],
);
if let Err(e) = done {
log::debug!("recording sidecar {path}: {e}");
}
}
/// Apply one sidecar's judgement to the image or images it describes.
///
/// # Why more than one image
///
/// A sidecar is named for the stem it shares with its photograph, so a RAW and
/// the JPEG the camera wrote beside it — one photograph under FR-CAT-11 — share
/// a document, and both rows have to carry the judgement or the grid disagrees
/// with itself depending on which of the pair it is showing.
///
/// # Why the match is verified in Rust
///
/// The query narrows the search to rows sharing the stem, and it is only a
/// filter: `library::sidecar_path` is what actually decides, applied to each
/// candidate. A judgement landing on the wrong photograph is a silent,
/// permanent wrong, so the rule lives in one place and the SQL only has to
/// return a superset of what it accepts.
///
/// Returns how many images gained a judgement they did not have.
pub(super) fn apply_judgement(
conn: &rusqlite::Connection,
root_id: i64,
sidecar: &str,
rating: u8,
flag: u8,
label: u8,
) -> Result<usize, String> {
let stem = sidecar
.rsplit_once('.')
.map(|(s, _)| s)
.unwrap_or(sidecar)
.to_string();
// Every name that starts `{stem}.`, as a range over the
// `(root_id, source_ref)` key: `/` is the byte after `.`, so the half-open
// range holds exactly those names. It was a `LIKE`, and a `LIKE` is
// case-insensitive, which no index here can serve -- every sidecar a pull
// took in read all 24,000 names of the root, 1.5 ms each. The names it
// matched beyond these differed only in case, and the check below has
// always refused them: `sidecar_path` compares exactly.
let candidates: Vec<(i64, String)> = {
let mut stmt = conn
.prepare_cached(
"SELECT id, source_ref FROM images
WHERE root_id = ?1 AND source_ref >= ?2 AND source_ref < ?3",
)
.map_err(|e| e.to_string())?;
let rows = stmt
.query_map(
rusqlite::params![root_id, format!("{stem}."), format!("{stem}/")],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.map_err(|e| e.to_string())?;
rows.filter_map(Result::ok)
.filter(|(_, source): &(i64, String)| sidecar_path(source) == sidecar)
.collect()
};
let mut applied = 0usize;
for (image, _) in candidates {
let id = dr_types::ImageId(image as u64);
let version =
dr_catalog::rating::default_version_id(conn, id).map_err(|e| e.to_string())?;
// The sidecar's value is taken, not the larger of the two. It is the
// authoritative store and the fuse above has already resolved any
// contest between devices on `revision` — so lowering a rating from
// four to one on the tablet has to lower it here, and taking a maximum
// would quietly refuse every demotion the photographer ever made.
//
// A zero is the one thing not taken: it means *never judged* rather
// than "judged zero", so a sidecar that carries none cannot erase a
// star this device holds. The same asymmetry, and the same direction
// of caution, as `dr_pipeline::sidecar::merge_judgement`. The cost is
// the one that rule always carries — clearing a rating does not
// propagate.
//
// TRACES: FR-CAT-5
// The label under the same rule: a sidecar with none leaves this
// device's label alone.
let changed = conn
.execute(
"UPDATE versions
SET rating = CASE WHEN ?2 > 0 THEN ?2 ELSE rating END,
flag = CASE WHEN ?3 > 0 THEN ?3 ELSE flag END,
label = CASE WHEN ?4 > 0 THEN ?4 ELSE label END
WHERE id = ?1
AND ((?2 > 0 AND rating <> ?2) OR (?3 > 0 AND flag <> ?3)
OR (?4 > 0 AND label IS NOT ?4))",
rusqlite::params![
version,
rating.min(5) as i64,
flag.min(2) as i64,
if label <= 5 { label as i64 } else { 0 }
],
)
.map_err(|e| e.to_string())?;
applied += changed;
}
Ok(applied)
}
/// Read back the folder ETags stored by a previous scan.
///
/// A failure here is not fatal — an empty map simply means no pruning, which
/// is correct but slower. Refusing to scan because the last scan's bookkeeping
/// is unreadable would be the worse outcome.
pub(super) fn load_folder_etags(
catalog: &Catalog,
root: &str,
) -> std::collections::HashMap<RemotePath, dr_sync::Validator> {
let mut out = std::collections::HashMap::new();
let sql = "SELECT f.path, f.etag FROM folders f
JOIN roots r ON r.id = f.root_id
WHERE r.label = ?1 AND f.etag IS NOT NULL";
let Ok(mut stmt) = catalog.connection().prepare(sql) else {
return out;
};
let rows = stmt.query_map([root], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
});
if let Ok(rows) = rows {
for (path, etag) in rows.flatten() {
out.insert(RemotePath::new(path), dr_sync::Validator::new(etag));
}
}
out
}
/// Write a scan's findings into the catalog.
///
/// Images insert at `metadata_state = 1` (stat-only): the scan knows name and
/// size but has read no EXIF, and pretending otherwise would make a date
/// filter silently wrong.
///
/// No job is queued. This used to enqueue a `Thumbnail` job per photograph,
/// and nothing has ever claimed that kind: the grid's worker and the thumbnail
/// sweep both find their work by asking the store what it lacks, which is the
/// only place that knows another device already made one. The reference
/// catalog held 23,582 of them, one per image, re-coalesced on every scan
/// (#73; catalog.md §6.1).
pub(super) fn persist(
catalog: &Catalog,
root: &str,
result: &dr_sync::ScanResult,
) -> Result<(), dr_catalog::CatalogError> {
let conn = catalog.connection();
let tx = conn.unchecked_transaction()?;
// One root row per library folder, reused across scans.
tx.execute(
"INSERT INTO roots(kind, label, last_seen) VALUES ('remote', ?1, ?2)
ON CONFLICT DO NOTHING",
rusqlite::params![root, now_secs()],
)?;
let root_id: i64 = tx.query_row(
"SELECT id FROM roots WHERE label = ?1 AND kind = 'remote'",
[root],
|r| r.get(0),
)?;
// Folder ETags first — without these persisted, the next scan prunes
// nothing and walks the whole tree again (ARCH §6.6).
{
let mut folder = tx.prepare_cached(
"INSERT INTO folders(root_id, path, etag) VALUES (?1, ?2, ?3)
ON CONFLICT(root_id, path) DO UPDATE SET etag = excluded.etag",
)?;
for (path, validator) in &result.directories {
folder.execute(rusqlite::params![
root_id,
path.as_str(),
validator.as_str()
])?;
}
}
// Every statement below runs once per photograph listed -- 1,600 for one
// folder whose sidecar changed, 24,000 for a first scan -- so each is
// prepared once, and a folder's id is looked up once per folder rather
// than once per photograph in it.
let mut folder_of =
tx.prepare_cached("SELECT id FROM folders WHERE root_id = ?1 AND path = ?2")?;
let mut folder_ids: std::collections::HashMap<String, Option<i64>> =
std::collections::HashMap::new();
// TRACES: FR-PLAT-AND-2 | FR-CAT-9
// The `availability` arm is what ends an offline library, and it does
// it one photograph at a time. 3 is `Availability::Offline` and 0 is
// `MetadataOnly`, the same code this statement inserts new rows with —
// so a row that was marked offline when the root became unreachable is
// returned to exactly the state a fresh scan would have given it, and
// a row that was never marked is not touched at all.
//
// Conditional rather than a blanket reset for the same reason
// `dr_catalog::walk` restores per file rather than per root: the only
// thing that may clear "I could not reach this" is having reached it,
// and this statement runs precisely once per file the scan listed.
let mut upsert = tx.prepare_cached(
"INSERT INTO images(root_id, folder_id, source_ref, format, file_size,
availability, metadata_state, added_at)
VALUES (?1, ?2, ?3, ?4, ?5, 0, 1, ?6)
ON CONFLICT(root_id, source_ref) DO UPDATE SET
file_size = excluded.file_size,
folder_id = excluded.folder_id,
availability = CASE WHEN images.availability = 3
THEN 0 ELSE images.availability END",
)?;
let mut image_of =
tx.prepare_cached("SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2")?;
let mut remote = tx.prepare_cached(
"INSERT INTO remote(image_id, file_id, etag, remote_path)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(image_id) DO UPDATE SET
etag = excluded.etag, remote_path = excluded.remote_path",
)?;
for entry in &result.images {
let folder_id: Option<i64> = match entry.path.parent() {
None => None,
Some(p) => match folder_ids.get(p.as_str()) {
Some(id) => *id,
None => {
let id = folder_of
.query_row(rusqlite::params![root_id, p.as_str()], |r| r.get(0))
.ok();
folder_ids.insert(p.as_str().to_string(), id);
id
}
},
};
upsert.execute(rusqlite::params![
root_id,
folder_id,
entry.path.as_str(),
entry
.path
.name()
.rsplit_once('.')
.map(|(_, e)| e.to_ascii_lowercase()),
entry.size as i64,
now_secs(),
])?;
let image_id: i64 = image_of
.query_row(rusqlite::params![root_id, entry.path.as_str()], |r| {
r.get(0)
})?;
// Remote identity, keyed on oc:fileid so a server-side move is a move
// rather than a re-download (FR-NC-5).
if let RemoteId::Stable(file_id) = entry.id {
remote.execute(rusqlite::params![
image_id,
file_id as i64,
entry.validator.as_str(),
entry.path.as_str()
])?;
}
}
drop((folder_of, upsert, image_of, remote));
tx.commit()?;
Ok(())
}
pub fn now_secs() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::library::test_support::*;
#[test]
fn persist_inserts_images_and_folder_etags() {
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![entry("PhotosRaw/a.CR2", 1001, 30_000_000)],
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1"))],
progress: Default::default(),
sidecars: Vec::new(),
};
persist(&catalog, "PhotosRaw", &result).unwrap();
assert_eq!(total_images(&catalog).unwrap(), 1);
// Folder ETags must persist or the next scan prunes nothing.
let etag: String = catalog
.connection()
.query_row(
"SELECT etag FROM folders WHERE path = 'PhotosRaw'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etag, "e1");
}
#[test]
fn images_land_as_stat_only_not_full_metadata() {
// The scan read no EXIF. Claiming otherwise would make a date filter
// silently wrong on a freshly scanned library.
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![entry("PhotosRaw/a.CR2", 1001, 30_000_000)],
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1"))],
progress: Default::default(),
sidecars: Vec::new(),
};
persist(&catalog, "PhotosRaw", &result).unwrap();
let state: i64 = catalog
.connection()
.query_row("SELECT metadata_state FROM images", [], |r| r.get(0))
.unwrap();
assert_eq!(state, 1);
}
#[test]
fn rescanning_updates_rather_than_duplicating() {
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![entry("PhotosRaw/a.CR2", 1001, 30_000_000)],
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1"))],
progress: Default::default(),
sidecars: Vec::new(),
};
persist(&catalog, "PhotosRaw", &result).unwrap();
persist(&catalog, "PhotosRaw", &result).unwrap();
assert_eq!(total_images(&catalog).unwrap(), 1, "no duplicate rows");
let roots: i64 = catalog
.connection()
.query_row("SELECT count(*) FROM roots", [], |r| r.get(0))
.unwrap();
assert_eq!(roots, 1, "no duplicate roots");
}
#[test]
fn stable_file_ids_are_recorded() {
// FR-NC-5: a server-side move must be a move, not a re-download.
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![entry("PhotosRaw/a.CR2", 4242, 30_000_000)],
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1"))],
progress: Default::default(),
sidecars: Vec::new(),
};
persist(&catalog, "PhotosRaw", &result).unwrap();
let file_id: i64 = catalog
.connection()
.query_row("SELECT file_id FROM remote", [], |r| r.get(0))
.unwrap();
assert_eq!(file_id, 4242);
}
#[test]
fn a_scan_queues_no_work_nothing_would_claim() {
// Every photograph used to leave a `Thumbnail` job behind, and no
// handler for that kind exists: thumbnails are owed by the store, not
// the queue. Scanned twice, because the second scan is the one that
// used to re-coalesce every row.
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![
entry("PhotosRaw/a.CR2", 1, 30_000_000),
entry("PhotosRaw/b.CR2", 2, 30_000_000),
],
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1"))],
progress: Default::default(),
sidecars: Vec::new(),
};
persist(&catalog, "PhotosRaw", &result).unwrap();
persist(&catalog, "PhotosRaw", &result).unwrap();
let jobs: i64 = catalog
.connection()
.query_row("SELECT count(*) FROM jobs", [], |r| r.get(0))
.unwrap();
assert_eq!(jobs, 0);
let images: i64 = catalog
.connection()
.query_row("SELECT count(*) FROM images", [], |r| r.get(0))
.unwrap();
assert_eq!(images, 2, "the scan still records what it found");
}
#[test]
fn folder_etags_round_trip_for_the_next_scan() {
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![],
directories: vec![
(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1")),
(
RemotePath::new("PhotosRaw/2026"),
dr_sync::Validator::new("e2"),
),
],
progress: Default::default(),
sidecars: Vec::new(),
};
persist(&catalog, "PhotosRaw", &result).unwrap();
let known = load_folder_etags(&catalog, "PhotosRaw");
assert_eq!(known.len(), 2);
assert_eq!(
known
.get(&RemotePath::new("PhotosRaw/2026"))
.map(|v| v.as_str()),
Some("e2")
);
}
/// A remote library holding the given files, indexed as a scan leaves it.
fn library(files: &[&str]) -> Catalog {
let cat = Catalog::in_memory().unwrap();
let c = cat.connection();
c.execute(
"INSERT INTO roots(id, kind, label) VALUES (1, 'remote', 'lib')",
[],
)
.unwrap();
for (i, f) in files.iter().enumerate() {
c.execute(
"INSERT INTO images(root_id, source_ref, added_at) VALUES (1, ?1, 0)",
[f],
)
.unwrap();
let image = c.last_insert_rowid();
c.execute(
"INSERT INTO remote(image_id, file_id) VALUES (?1, ?2)",
rusqlite::params![image, 1000 + i as i64],
)
.unwrap();
}
dr_catalog::rating::ensure_default_versions(c).unwrap();
cat
}
fn judgements(cat: &Catalog) -> Vec<(String, i64, i64)> {
let c = cat.connection();
let mut stmt = c
.prepare(
"SELECT i.source_ref, v.rating, v.flag FROM images i
JOIN versions v ON v.image_id = i.id AND v.is_default = 1
ORDER BY i.source_ref",
)
.unwrap();
let v = stmt
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))
.unwrap()
.map(Result::unwrap)
.collect();
v
}
#[test]
fn a_label_from_another_device_reaches_the_catalog_and_none_clears_nothing() {
// TRACES: FR-CAT-5
let cat = library(&["2026/a.CR2"]);
assert_eq!(
apply_judgement(cat.connection(), 1, "2026/a.drsc", 0, 0, 4).unwrap(),
1
);
let label = |cat: &Catalog| -> Option<i64> {
cat.connection()
.query_row("SELECT label FROM versions WHERE is_default = 1", [], |r| {
r.get(0)
})
.unwrap()
};
assert_eq!(label(&cat), Some(4));
apply_judgement(cat.connection(), 1, "2026/a.drsc", 3, 0, 0).unwrap();
assert_eq!(
label(&cat),
Some(4),
"a sidecar with no label left it alone"
);
}
/// The regression, in one line: a rating in a sidecar reaches the grid.
/// Before this there was no path by which it could.
#[test]
fn a_rating_from_another_device_reaches_the_catalog() {
let cat = library(&["2026/a.CR2"]);
let n = apply_judgement(cat.connection(), 1, "2026/a.drsc", 4, 1, 0).unwrap();
assert_eq!(n, 1);
assert_eq!(judgements(&cat), vec![("2026/a.CR2".to_string(), 4, 1)]);
}
/// A RAW and the JPEG beside it are one photograph (FR-CAT-11) sharing one
/// sidecar, so both rows have to carry the judgement — otherwise the grid
/// disagrees with itself depending which of the pair it draws.
#[test]
fn both_halves_of_a_raw_and_jpeg_pair_are_judged() {
let cat = library(&["2026/a.CR2", "2026/a.jpg"]);
let n = apply_judgement(cat.connection(), 1, "2026/a.drsc", 3, 0, 0).unwrap();
assert_eq!(n, 2);
assert_eq!(
judgements(&cat),
vec![
("2026/a.CR2".to_string(), 3, 0),
("2026/a.jpg".to_string(), 3, 0)
]
);
}
/// The lookup is a range over names starting `{stem}.`, so the names
/// either side of that range — a longer stem, a case variant, a folder
/// whose name holds the wildcards the old `LIKE` had to escape — must
/// neither be missed nor caught.
#[test]
fn a_sidecar_reaches_exactly_its_own_photographs() {
let cat = library(&[
"2026/50%_off/a.CR2",
"2026/50%_off/a.jpg",
"2026/50%_off/A.CR2",
"2026/50%_off/a.b.CR2",
"2026/50%_off/ab.CR2",
"2026/50%_off/a/b.CR2",
"2026/50Xxoff/a.CR2",
]);
let n = apply_judgement(cat.connection(), 1, "2026/50%_off/a.drsc", 2, 0, 0).unwrap();
assert_eq!(n, 2);
let judged: Vec<String> = judgements(&cat)
.into_iter()
.filter(|(_, rating, _)| *rating == 2)
.map(|(path, _, _)| path)
.collect();
assert_eq!(judged, vec!["2026/50%_off/a.CR2", "2026/50%_off/a.jpg"]);
}
/// A demotion has to travel. Taking the larger of the two would refuse
/// every rating the photographer ever lowered — and lowering one is most
/// of what a second pass over a shoot does.
#[test]
fn a_lowered_rating_travels() {
let cat = library(&["2026/a.CR2"]);
apply_judgement(cat.connection(), 1, "2026/a.drsc", 4, 0, 0).unwrap();
apply_judgement(cat.connection(), 1, "2026/a.drsc", 1, 0, 0).unwrap();
assert_eq!(judgements(&cat)[0].1, 1);
}
/// But a zero is *unjudged*, not "judged zero". A device that never culled
/// the frame must not erase the stars of one that did.
#[test]
fn an_unjudged_sidecar_does_not_erase_a_local_rating() {
let cat = library(&["2026/a.CR2"]);
apply_judgement(cat.connection(), 1, "2026/a.drsc", 5, 2, 0).unwrap();
assert_eq!(
apply_judgement(cat.connection(), 1, "2026/a.drsc", 0, 0, 0).unwrap(),
0
);
assert_eq!(judgements(&cat), vec![("2026/a.CR2".to_string(), 5, 2)]);
}
/// Idempotent: a scan runs repeatedly, and a sidecar whose judgement is
/// already in the catalog must report no change or the status line claims
/// work that did not happen.
#[test]
fn applying_the_same_judgement_twice_changes_nothing() {
let cat = library(&["2026/a.CR2"]);
assert_eq!(
apply_judgement(cat.connection(), 1, "2026/a.drsc", 4, 1, 0).unwrap(),
1
);
assert_eq!(
apply_judgement(cat.connection(), 1, "2026/a.drsc", 4, 1, 0).unwrap(),
0
);
}
/// The `LIKE` is a filter and not the decision. A stem holding a wildcard
/// would otherwise reach photographs it has nothing to do with, and a
/// rating landing on the wrong frame is silent and permanent.
#[test]
fn a_wildcard_in_a_path_does_not_reach_another_photograph() {
// `_` is LIKE's single-character wildcard, so an unescaped `a_b` stem
// would also match `axb`.
let cat = library(&["2026/a_b.CR2", "2026/axb.CR2"]);
apply_judgement(cat.connection(), 1, "2026/a_b.drsc", 5, 0, 0).unwrap();
assert_eq!(
judgements(&cat),
vec![
("2026/a_b.CR2".to_string(), 5, 0),
("2026/axb.CR2".to_string(), 0, 0)
]
);
}
/// A stem that is a prefix of another must not spill onto it: `a.drsc`
/// describes `a.CR2`, never `ab.CR2`.
#[test]
fn a_shared_prefix_is_not_a_shared_sidecar() {
let cat = library(&["2026/a.CR2", "2026/ab.CR2"]);
apply_judgement(cat.connection(), 1, "2026/a.drsc", 5, 0, 0).unwrap();
assert_eq!(
judgements(&cat),
vec![
("2026/a.CR2".to_string(), 5, 0),
("2026/ab.CR2".to_string(), 0, 0)
]
);
}
/// A sidecar for a photograph this device has not indexed is not an error:
/// the scan may have pruned the folder, or the file may be a format this
/// device does not accept.
#[test]
fn a_sidecar_with_no_photograph_here_is_not_a_failure() {
let cat = library(&["2026/a.CR2"]);
assert_eq!(
apply_judgement(cat.connection(), 1, "2026/elsewhere.drsc", 5, 0, 0).unwrap(),
0
);
}
/// What a scan pass and a sidecar pull cost on a real library.
///
/// DR_BENCH_CATALOG=/path/to/copy.sqlite cargo test --release -p dr-ui \
/// --lib persist_bench -- --ignored --nocapture --test-threads=1
///
/// Replays the catalog's own `remote` rows through [`persist`] — the
/// largest folder, which is what a pass relists after one sidecar write
/// there, and then the whole library, which is a first scan — and looks
/// up every `.drsc` sidecar the catalog has read, as [`apply_judgement`]
/// does on a pull (with no judgement, so nothing is written). Works on a
/// scratch copy beside the catalog it is given, and prints a fingerprint
/// of what `persist` left so two builds can be shown to agree.
#[test]
#[ignore]
fn persist_bench() {
use std::time::{Duration, Instant};
let Ok(source) = std::env::var("DR_BENCH_CATALOG") else {
eprintln!("DR_BENCH_CATALOG not set; nothing to measure");
return;
};
let scratch = format!("{source}.persist-bench");
let _ = std::fs::remove_file(&scratch);
{
let from = rusqlite::Connection::open(&source).unwrap();
from.execute("VACUUM INTO ?1", [&scratch]).unwrap();
}
let catalog = Catalog::open(std::path::Path::new(&scratch)).unwrap();
let conn = catalog.connection();
let root: String = conn
.query_row(
"SELECT label FROM roots WHERE kind = 'remote' LIMIT 1",
[],
|r| r.get(0),
)
.unwrap();
let entries = |filter: &str| -> dr_sync::ScanResult {
let mut q = conn
.prepare(&format!(
"SELECT i.source_ref, r.file_id, r.etag, coalesce(i.file_size, 0)
FROM remote r JOIN images i ON i.id = r.image_id
WHERE 1 {filter}
ORDER BY i.source_ref"
))
.unwrap();
let images: Vec<dr_sync::RemoteEntry> = q
.query_map([], |r| {
Ok(dr_sync::RemoteEntry {
id: RemoteId::Stable(r.get::<_, i64>(1)? as u64),
path: RemotePath::new(r.get::<_, String>(0)?),
kind: dr_sync::EntryKind::File,
validator: dr_sync::Validator::new(
r.get::<_, Option<String>>(2)?.unwrap_or_default(),
),
size: r.get::<_, i64>(3)? as u64,
modified: None,
has_preview: false,
materialised: true,
})
})
.unwrap()
.collect::<Result<_, _>>()
.unwrap();
let mut dirs: Vec<String> = images
.iter()
.filter_map(|e| e.path.parent().map(|p| p.as_str().to_string()))
.collect();
dirs.sort();
dirs.dedup();
let directories = dirs
.into_iter()
.map(|d| {
let etag: Option<String> = conn
.query_row("SELECT etag FROM folders WHERE path = ?1", [&d], |r| {
r.get(0)
})
.ok()
.flatten();
(
RemotePath::new(d),
dr_sync::Validator::new(etag.unwrap_or_default()),
)
})
.collect();
dr_sync::ScanResult {
images,
directories,
progress: Default::default(),
sidecars: Vec::new(),
}
};
let biggest: i64 = conn
.query_row(
"SELECT folder_id FROM images GROUP BY folder_id ORDER BY count(*) DESC LIMIT 1",
[],
|r| r.get(0),
)
.unwrap();
let folder = entries(&format!("AND i.folder_id = {biggest}"));
let all = entries("");
let cpu_now = || {
// The test's own thread: the harness runs it off the main one.
std::fs::read_to_string("/proc/thread-self/schedstat")
.ok()
.and_then(|s| s.split_whitespace().next()?.parse::<u64>().ok())
.map(Duration::from_nanos)
.unwrap_or_default()
};
let time = |label: &str, runs: usize, f: &mut dyn FnMut()| {
let mut wall = Vec::new();
let mut cpu = Vec::new();
for _ in 0..runs {
let (c, t) = (cpu_now(), Instant::now());
f();
wall.push(t.elapsed());
cpu.push(cpu_now().saturating_sub(c));
}
wall.sort();
cpu.sort();
println!(
"{label:44} best {:9.2} ms median {:9.2} ms cpu {:9.2} ms",
wall[0].as_secs_f64() * 1e3,
wall[runs / 2].as_secs_f64() * 1e3,
cpu[0].as_secs_f64() * 1e3
);
};
time(
&format!("persist, largest folder ({} images)", folder.images.len()),
5,
&mut || persist(&catalog, &root, &folder).unwrap(),
);
time(
&format!("persist, whole library ({} images)", all.images.len()),
3,
&mut || persist(&catalog, &root, &all).unwrap(),
);
let root_id: i64 = conn
.query_row("SELECT id FROM roots WHERE label = ?1", [&root], |r| {
r.get(0)
})
.unwrap();
let sidecars: Vec<String> = {
let mut q = conn
.prepare("SELECT path FROM sidecars WHERE path LIKE '%.drsc' ORDER BY path")
.unwrap();
q.query_map([], |r| r.get(0))
.unwrap()
.collect::<Result<_, _>>()
.unwrap()
};
let mut matched = 0usize;
time(
&format!("apply_judgement lookup x{}", sidecars.len()),
3,
&mut || {
for s in &sidecars {
matched += apply_judgement(conn, root_id, s, 0, 0, 0).unwrap();
}
},
);
let fingerprint: (i64, i64, i64, i64) = conn
.query_row(
"SELECT (SELECT count(*) || ':' || total(id * 31 + coalesce(folder_id, 0) * 7
+ file_size % 1000003 + availability * 3
+ metadata_state + length(source_ref))
FROM images),
(SELECT total(image_id * 13 + file_id % 1000003 + length(etag)
+ length(remote_path)) FROM remote),
(SELECT count(*) || ':' || total(kind * 17 + coalesce(subject_id, 0) * 5
+ priority * 3 + state + attempts) FROM jobs),
(SELECT count(*) || ':' || total(length(path) + length(etag)) FROM folders)",
[],
|r| {
let h = |s: String| {
s.bytes()
.fold(0i64, |a, b| a.wrapping_mul(131).wrapping_add(b as i64))
};
Ok((
h(r.get::<_, String>(0)?),
h(r.get::<_, f64>(1)?.to_string()),
h(r.get::<_, String>(2)?),
h(r.get::<_, String>(3)?),
))
},
)
.unwrap();
println!("fingerprint after persist: {fingerprint:?}, judgements matched {matched}");
drop(catalog);
let _ = std::fs::remove_file(&scratch);
let _ = std::fs::remove_file(format!("{scratch}-wal"));
let _ = std::fs::remove_file(format!("{scratch}-shm"));
}
}