Thirty-nine spawn sites in dr-ui, and one in the Android entry point, called std::thread::spawn or a Builder of their own, and most of the threads they started were <unnamed> in a panic message or a profiler. Each now calls executors::spawn with its executor and a role, so the thread is named <executor>:<role> — net:sync, decode:thumbs, io:catalog-open — and knows which executor it is on. The three that already set a name (automation, import, prefetch) keep their name as the role. Behaviour is unchanged: each job still gets a thread of its own when it starts, and spawn panics where std::thread::spawn did. The module's documentation now says how a job is assigned: by what it spends its time on, so a sweep that fetches bytes and then decodes them is Decode, and a sidecar write that touches the catalog is Network. Left as they were: the segmentation and refine workers in masks_ui.rs, which another change is reworking, and test-only threads.
1279 lines
49 KiB
Rust
1279 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 crate::executors::{self, Executor};
|
|
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();
|
|
|
|
executors::spawn(Executor::Network, "scan", 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"));
|
|
}
|
|
}
|