Files
DarkRoom/ui/dr-ui/src/library/scan.rs
T
dtourolle 454375243c Measure what opening the catalog, a sync pass and a scan cost on a real library
Two benches for reading side by side before and after a change, against a
copy of a real catalog, in the manner of identity_bench:

- `dr-catalog --example catalog_bench CATALOG [FACES_DIR]` times
  `Catalog::open` and the backfill inside it step by step, the upload
  snapshot, a merge of the catalog with a copy of itself, and the face
  shard export and import in the steady state where nothing is new.

- `persist_bench`, an ignored test in dr-ui's scan module because
  `persist` and `apply_judgement` are private to it, replays the
  catalog's own rows through `persist` (the largest folder, and the whole
  library) and looks up every `.drsc` sidecar the catalog has read. It
  works on a scratch copy and prints a fingerprint of what `persist` left,
  so two builds can be shown to agree.

Both print best, median and CPU time; the CPU figure is the one to compare
while other builds share the machine.
2026-09-25 22:06:58 -04:00

1251 lines
47 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, JobKind, Priority};
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 `LIKE` 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 path holding a `%`, a `_` or a bracket would otherwise match
/// more than it should, and a judgement landing on the wrong photograph is a
/// silent, permanent wrong.
///
/// 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();
// The escape is the point: `%` and `_` are wildcards, and a photographer's
// folder is entitled to contain both.
let prefix = stem
.replace('\\', "\\\\")
.replace('%', "\\%")
.replace('_', "\\_");
let candidates: Vec<(i64, String)> = {
let mut stmt = conn
.prepare(
"SELECT id, source_ref FROM images
WHERE root_id = ?1 AND source_ref LIKE ?2 ESCAPE '\\'",
)
.map_err(|e| e.to_string())?;
let rows = stmt
.query_map(rusqlite::params![root_id, format!("{prefix}.%")], |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. A `Thumbnail` job is enqueued per image, coalescing
/// with anything already pending.
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).
for (path, validator) in &result.directories {
tx.execute(
"INSERT INTO folders(root_id, path, etag) VALUES (?1, ?2, ?3)
ON CONFLICT(root_id, path) DO UPDATE SET etag = excluded.etag",
rusqlite::params![root_id, path.as_str(), validator.as_str()],
)?;
}
for entry in &result.images {
let folder_id: Option<i64> = entry.path.parent().and_then(|p| {
tx.query_row(
"SELECT id FROM folders WHERE root_id = ?1 AND path = ?2",
rusqlite::params![root_id, p.as_str()],
|r| r.get(0),
)
.ok()
});
// 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.
tx.execute(
"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",
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 = tx.query_row(
"SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2",
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 {
tx.execute(
"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",
rusqlite::params![
image_id,
file_id as i64,
entry.validator.as_str(),
entry.path.as_str()
],
)?;
}
}
tx.commit()?;
// Thumbnail jobs after the commit, so a failure mid-insert does not leave
// jobs pointing at rows that never landed.
for entry in &result.images {
if let Ok(image_id) = conn.query_row(
"SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2",
rusqlite::params![root_id, entry.path.as_str()],
|r| r.get::<_, i64>(0),
) {
let _ = dr_catalog::jobs::enqueue(
conn,
JobKind::Thumbnail,
Some(image_id),
Priority::Background,
None,
);
}
}
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_thumbnail_job_is_queued_per_image() {
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();
let jobs: i64 = catalog
.connection()
.query_row("SELECT count(*) FROM jobs WHERE kind = 2", [], |r| r.get(0))
.unwrap();
assert_eq!(jobs, 2);
}
#[test]
fn a_second_scan_does_not_multiply_jobs() {
let catalog = Catalog::in_memory().unwrap();
let result = dr_sync::ScanResult {
images: vec![entry("PhotosRaw/a.CR2", 1, 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 WHERE kind = 2", [], |r| r.get(0))
.unwrap();
assert_eq!(jobs, 1, "coalesced, not queued twice");
}
#[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)
]
);
}
/// 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"));
}
}