`write_one_sidecar` took an `Option<&NextcloudBackend>` it ignored. It was a leftover from the shape the write path had before online and offline separated into two functions: the offline one records to the cache and queues, and has no server to talk to by definition. A parameter that is always `None` and always unused says the opposite — that there is a case where it is `Some` — and the next reader has to check. The closure it was threaded through is renamed to say what it does rather than how it is called. `run(None)` needed the reader to know what `None` meant; `queue_all()` is the sentence. Also regenerates the traceability matrix, which now records FR-CAT-9's queue. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
3452 lines
129 KiB
Rust
3452 lines
129 KiB
Rust
//! TRACES: FR-CAT-1 | FR-CAT-4 | FR-NC-3 | NFR-P9
|
||
//! Opening a remote library: scan → catalog → grid.
|
||
//!
|
||
//! This is the wire between three pieces that already worked separately —
|
||
//! `dr_sync::scan` walks the tree, `dr_catalog` indexes it, and
|
||
//! `dr_decode::preview` turns bytes into pixels. Until now the "Open library"
|
||
//! button logged its intent and stopped.
|
||
//!
|
||
//! # Threading
|
||
//!
|
||
//! Slint's event loop is single-threaded and must never block (NFR-P9), so
|
||
//! every network and decode operation runs on a worker thread and results
|
||
//! return through an mpsc channel drained by a Slint timer. That is the same
|
||
//! shape the login flow uses; it is repeated rather than shared because the
|
||
//! message types differ and a generic version would obscure both.
|
||
//!
|
||
//! # Why thumbnails are fetched, not derived from the scan
|
||
//!
|
||
//! A scan yields paths and sizes, nothing visual. Each thumbnail costs its own
|
||
//! range request, so they are fetched **only for cells the grid actually
|
||
//! wants** — never for the whole library up front. On the reference library
|
||
//! that is the difference between a few MB and ~370 GB (ARCH §6.7).
|
||
|
||
use std::path::PathBuf;
|
||
use std::sync::mpsc::{Receiver, Sender};
|
||
|
||
use dr_catalog::{Catalog, JobKind, Priority};
|
||
use dr_sync::{RemoteBackend, RemoteId, RemotePath};
|
||
use dr_sync_nextcloud::{AppCredentials, NextcloudBackend};
|
||
use dr_thumbs::ThumbStore;
|
||
|
||
use crate::sidecar_cache::SidecarCache;
|
||
use dr_types::FormatFilter;
|
||
|
||
/// Largest preview worth fetching whole.
|
||
///
|
||
/// A located preview above this is skipped rather than transferred: past a few
|
||
/// MB the saving over the full file stops justifying the wait, and a 256px
|
||
/// thumbnail needs nothing like that much detail.
|
||
const MAX_PREVIEW_BYTES: u64 = 8 * 1024 * 1024;
|
||
|
||
/// 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,
|
||
},
|
||
/// 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.
|
||
Failed { message: String, offline: bool },
|
||
}
|
||
|
||
/// One decoded thumbnail, ready for the grid.
|
||
#[derive(Debug)]
|
||
pub struct ThumbnailReady {
|
||
/// Index into the grid model this belongs to.
|
||
pub row: usize,
|
||
pub width: u32,
|
||
pub height: u32,
|
||
pub rgba: Vec<u8>,
|
||
/// Whether these pixels came off local disk rather than the server.
|
||
///
|
||
/// The grid paints both identically, so this exists solely for
|
||
/// reachability: a store hit is evidence about the *cache*, not the
|
||
/// network, and treating one as proof of connectivity clears offline mode
|
||
/// before a single request has been attempted.
|
||
pub from_cache: bool,
|
||
}
|
||
|
||
/// Capture metadata read from the same header the thumbnail needed.
|
||
///
|
||
/// Free: the header fetch happens either way, so parsing EXIF out of it costs
|
||
/// no extra transfer. That is what fills the timeline as the user browses,
|
||
/// rather than a separate 6 GB sweep over the library.
|
||
#[derive(Debug, Clone)]
|
||
pub struct MetadataFound {
|
||
pub image_id: i64,
|
||
pub captured_at: Option<i64>,
|
||
pub captured_offset: Option<i32>,
|
||
pub camera: Option<String>,
|
||
pub lens: Option<String>,
|
||
pub iso: Option<u32>,
|
||
}
|
||
|
||
/// Messages from the thumbnail worker.
|
||
#[derive(Debug)]
|
||
pub enum ThumbnailMessage {
|
||
Ready(Box<ThumbnailReady>),
|
||
/// No preview could be extracted. The cell stays a placeholder rather than
|
||
/// silently retrying forever.
|
||
Unavailable {
|
||
row: usize,
|
||
reason: String,
|
||
},
|
||
/// How the batch split between the store and the network.
|
||
///
|
||
/// Sent once, before any fetch. Without it there is no way to tell a
|
||
/// working cache from a broken one — both fill the grid, one just costs
|
||
/// nothing.
|
||
Plan {
|
||
cached: usize,
|
||
fetching: usize,
|
||
dating: usize,
|
||
},
|
||
/// One header-only date read is starting.
|
||
///
|
||
/// Reported separately from thumbnail progress: this work produces no
|
||
/// visible cell, so without it the window looks idle while it runs.
|
||
DateProgress,
|
||
/// Capture dates were written to the catalog.
|
||
///
|
||
/// The timeline is rebuilt on this rather than per image — a histogram
|
||
/// that redrew 120 times during a batch would flicker for no benefit.
|
||
DatesRecorded(usize),
|
||
/// TRACES: FR-CAT-9
|
||
/// The server could not be reached while filling this batch.
|
||
///
|
||
/// Distinct from a run of [`Unavailable`](Self::Unavailable): those are
|
||
/// per-image verdicts ("this file has no extractable preview") and leave
|
||
/// the rest of the library alone, where this is a statement about the
|
||
/// connection. Sent at most once per batch, because a dropped connection
|
||
/// produces one of these per *cell* otherwise and the banner would be
|
||
/// rewritten sixty times.
|
||
Offline {
|
||
reason: String,
|
||
},
|
||
}
|
||
|
||
/// TRACES: FR-CAT-15 | FR-CAT-11
|
||
/// What it means for an image to be visible in the library.
|
||
///
|
||
/// Two exclusions, for two different reasons, and both must appear in *every*
|
||
/// query that counts or lists cells — the grid, the timeline, the metadata
|
||
/// sweep. A predicate present in four of five places is worse than absent: the
|
||
/// counts disagree with the cells and neither looks wrong on its own.
|
||
///
|
||
/// - `shadowed_by IS NULL` — a JPEG the camera wrote beside its RAW is that
|
||
/// same frame, not a second photograph.
|
||
/// - `trashed_at IS NULL` — a soft-deleted image has been moved to the trash
|
||
/// folder and is listed only by the trash view.
|
||
const VISIBLE: &str = "i.shadowed_by IS NULL AND i.trashed_at IS NULL";
|
||
|
||
/// [`VISIBLE`] for queries that do not alias `images`.
|
||
const VISIBLE_UNALIASED: &str = "shadowed_by IS NULL AND trashed_at IS NULL";
|
||
|
||
/// TRACES: FR-CAT-15
|
||
/// What the *trash view* lists: exactly what [`VISIBLE`] excludes on the second
|
||
/// clause, and still excludes on the first.
|
||
///
|
||
/// The inversion is deliberate and only correct on `trashed_at`. A shadowed JPEG
|
||
/// is not a separate photograph in the trash any more than it is in the library
|
||
/// — trashing a RAW takes its sibling with it, and listing both would offer to
|
||
/// restore the same frame twice.
|
||
const TRASHED: &str = "i.shadowed_by IS NULL AND i.trashed_at IS NOT NULL";
|
||
|
||
/// TRACES: FR-CAT-6 | FR-CULL-4
|
||
/// What the grid is narrowed to by the rating filter bar.
|
||
///
|
||
/// Applied in **SQL**, not by filtering the rows after reading them. On a
|
||
/// remote library a drawn-then-hidden cell has already cost a thumbnail
|
||
/// fetch, which is the transfer FR-NC-3 exists to avoid — and the count in the
|
||
/// header has to agree with the cells, which it cannot if the two are computed
|
||
/// at different stages.
|
||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
|
||
pub struct RatingFilter {
|
||
/// Minimum stars. 0 means no star constraint.
|
||
pub min_rating: u8,
|
||
/// Only images nothing has judged yet — neither starred nor flagged.
|
||
/// This is what lets a culling session resume where it stopped.
|
||
pub unjudged: bool,
|
||
/// `None` for no flag constraint, otherwise exactly that flag.
|
||
pub flag: Option<dr_types::FlagState>,
|
||
/// TRACES: FR-CAT-9
|
||
/// Only images whose original is stored on this device.
|
||
///
|
||
/// Carried here, beside the rating terms, because every query path already
|
||
/// threads this one struct: adding a parallel parameter to
|
||
/// `read_cells_scoped`, `read_cells_all` and both counts would give four
|
||
/// call sites the chance to disagree about what the grid is showing, and
|
||
/// the count disagreeing with the cells is the specific bug this type's
|
||
/// "filter in SQL" rule exists to prevent.
|
||
pub local_only: bool,
|
||
}
|
||
|
||
impl RatingFilter {
|
||
/// Whether this narrows anything, so the caller can skip the join.
|
||
pub fn is_unfiltered(&self) -> bool {
|
||
self.min_rating == 0 && !self.unjudged && self.flag.is_none() && !self.local_only
|
||
}
|
||
|
||
/// The SQL predicate, against an `images` aliased as `i`.
|
||
///
|
||
/// Returns a `String` of conditions ANDed together, or an empty string
|
||
/// where nothing is constrained. Every branch is built from integers this
|
||
/// code owns — no caller text reaches the SQL, so there is nothing to
|
||
/// escape.
|
||
///
|
||
/// A correlated subquery per term rather than a join to `versions`: an
|
||
/// image with no version row must still be *findable* as unrated, and an
|
||
/// inner join would silently drop exactly those images — the ones a
|
||
/// library scanned before ratings existed consists entirely of.
|
||
fn sql(&self) -> String {
|
||
let mut terms = Vec::new();
|
||
|
||
if self.min_rating > 0 {
|
||
terms.push(format!(
|
||
"coalesce((SELECT dv.rating FROM versions dv
|
||
WHERE dv.image_id = i.id AND dv.is_default = 1
|
||
LIMIT 1), 0) >= {}",
|
||
self.min_rating
|
||
));
|
||
}
|
||
|
||
if self.unjudged {
|
||
// Both axes: a frame that was picked but never starred has been
|
||
// judged, and re-presenting it would undo the user's decision to
|
||
// move past it.
|
||
terms.push(
|
||
"coalesce((SELECT dv.rating FROM versions dv
|
||
WHERE dv.image_id = i.id AND dv.is_default = 1
|
||
LIMIT 1), 0) = 0
|
||
AND coalesce((SELECT dv.flag FROM versions dv
|
||
WHERE dv.image_id = i.id AND dv.is_default = 1
|
||
LIMIT 1), 0) = 0"
|
||
.to_string(),
|
||
);
|
||
}
|
||
|
||
if let Some(flag) = self.flag {
|
||
terms.push(format!(
|
||
"coalesce((SELECT dv.flag FROM versions dv
|
||
WHERE dv.image_id = i.id AND dv.is_default = 1
|
||
LIMIT 1), 0) = {}",
|
||
flag_code(flag)
|
||
));
|
||
}
|
||
|
||
if self.local_only {
|
||
// `tier_actual`, not `tier_desired`: the question is what is
|
||
// *here*, not what a pin has promised will be. An image queued for
|
||
// download is exactly the one that cannot be opened yet, so
|
||
// showing it under "on this device" would be the wrong answer to
|
||
// the only question this filter is asked.
|
||
terms.push(format!(
|
||
"EXISTS (SELECT 1 FROM image_cache ic
|
||
WHERE ic.image_id = i.id AND ic.tier_actual >= {})",
|
||
dr_types::Tier::Original.stored()
|
||
));
|
||
}
|
||
|
||
if terms.is_empty() {
|
||
String::new()
|
||
} else {
|
||
format!(" AND ({})", terms.join(") AND ("))
|
||
}
|
||
}
|
||
}
|
||
|
||
/// The stored integer for a flag, matching `dr_catalog::rating`'s encoding.
|
||
fn flag_code(f: dr_types::FlagState) -> i64 {
|
||
match f {
|
||
dr_types::FlagState::Unflagged => 0,
|
||
dr_types::FlagState::Pick => 1,
|
||
dr_types::FlagState::Reject => 2,
|
||
}
|
||
}
|
||
|
||
/// TRACES: FR-CAT-8 | FR-NC-8 | FR-CULL-4 | FR-DEV-6
|
||
/// One amendment to one image's sidecar, on its way to the server.
|
||
#[derive(Debug, Clone)]
|
||
pub struct SidecarWrite {
|
||
/// Remote path of the *image*. The sidecar sits beside it, with the
|
||
/// extension replaced — that adjacency is what makes a sidecar findable
|
||
/// without an index (ARCH §6.12).
|
||
pub image_path: String,
|
||
pub version_uuid: String,
|
||
pub amendment: Amendment,
|
||
}
|
||
|
||
/// What a write changes about the version it names.
|
||
///
|
||
/// An enum rather than a struct of optional fields because the two are written
|
||
/// by different actions with different failure costs, and because a write must
|
||
/// never carry a *stale* copy of what it is not changing. A settings write that
|
||
/// also carried a rating would have to have read one from somewhere, and the
|
||
/// obvious somewhere — the catalog, moments earlier — is exactly how a cull
|
||
/// made between the read and the write gets silently reverted.
|
||
///
|
||
/// Everything not named by the variant is left as the file had it, which is
|
||
/// what makes the read-modify-write in [`write_one_sidecar`] a genuine
|
||
/// amendment rather than a replacement.
|
||
#[derive(Debug, Clone)]
|
||
pub enum Amendment {
|
||
/// A star rating and a pick/reject flag — the cull.
|
||
Judgement { rating: u8, flag: u8 },
|
||
/// TRACES: FR-DEV-6
|
||
/// Copied develop settings, applied within `scope`.
|
||
///
|
||
/// Carries the [`Scope`] rather than a pre-filtered preset so the target's
|
||
/// own framing can be spared *at the file*: excluding framing means
|
||
/// leaving the keys already in the sidecar untouched, which cannot be
|
||
/// expressed by the parameter list alone.
|
||
Settings {
|
||
preset: dr_pipeline::Preset,
|
||
scope: dr_pipeline::Scope,
|
||
},
|
||
}
|
||
|
||
/// Where an image's sidecar lives.
|
||
///
|
||
/// The image's own path with the extension replaced, not appended: `a.CR2`
|
||
/// becomes `a.drsc`, so a RAW and the JPEG beside it share one sidecar and
|
||
/// therefore one judgement. That is the intended behaviour — they are the same
|
||
/// photograph (FR-CAT-11), and the pairing logic in `dr_catalog::schema`
|
||
/// already treats them so.
|
||
pub fn sidecar_path(image_path: &str) -> String {
|
||
let stem = match image_path.rsplit_once('.') {
|
||
// Only an extension in the final segment counts; a dot in a directory
|
||
// name must not truncate the path.
|
||
Some((stem, ext)) if !ext.contains('/') => stem,
|
||
_ => image_path,
|
||
};
|
||
format!("{stem}.{}", dr_pipeline::sidecar::EXTENSION)
|
||
}
|
||
|
||
/// TRACES: FR-CAT-8 | FR-CAT-9 | FR-NC-10
|
||
/// Persist amendments to sidecars beside their images.
|
||
///
|
||
/// # Why this reads before it writes
|
||
///
|
||
/// A sidecar is the authoritative store and may already hold an edit made on
|
||
/// this or another device. Writing a fresh document containing only a rating
|
||
/// would delete that edit — the exact silent data loss the format's
|
||
/// unknown-key preservation exists to prevent. So each file is fetched,
|
||
/// parsed, amended, and written back; a fetch that 404s simply means there is
|
||
/// no sidecar yet and a new one is created.
|
||
///
|
||
/// # Why the local write is the commit point
|
||
///
|
||
/// FR-CAT-9 requires that edits made offline *queue and apply when the source
|
||
/// returns*. So every amendment is written to the local cache first and the
|
||
/// upload is best-effort: an entry stays marked pending until the server has
|
||
/// actually taken it, and [`spawn_outbox_drain`] retries the marked ones later.
|
||
///
|
||
/// This is what makes `offline` a parameter rather than a reason to skip. It
|
||
/// was one: a cull or a paste made with no connection used to be dropped
|
||
/// entirely, which for a pasted edit meant it survived nowhere at all — the
|
||
/// catalog holds no parameters. Now the two cases differ only in whether the
|
||
/// upload is attempted.
|
||
///
|
||
/// # Why failure here is logged rather than surfaced
|
||
///
|
||
/// The write has already succeeded locally by the time the network is touched,
|
||
/// so nothing is lost by a failure and there is nothing for the user to do
|
||
/// about it. Interrupting a cull with an error dialog per frame would be far
|
||
/// worse than the risk. The counts are reported once, at the end.
|
||
pub fn spawn_sidecar_writes(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
writes: Vec<SidecarWrite>,
|
||
cache_dir: PathBuf,
|
||
offline: bool,
|
||
) -> Receiver<SidecarMessage> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
let cache = SidecarCache::open(cache_dir);
|
||
|
||
// The runtime and the backend are only needed to *upload*. Offline,
|
||
// neither is built — and a failure to build either is not a failure to
|
||
// record the edit, it just means every write is queued instead.
|
||
let rt = if offline {
|
||
None
|
||
} else {
|
||
match crate::net_runtime::build() {
|
||
Ok(rt) => Some(rt),
|
||
Err(e) => {
|
||
log::debug!("no runtime for sidecar upload ({e}); queueing");
|
||
None
|
||
}
|
||
}
|
||
};
|
||
|
||
// Every write, recorded locally and queued. The offline path, and the
|
||
// fallback whenever a backend could not be built.
|
||
let queue_all = || {
|
||
let mut report = SidecarReport::default();
|
||
for w in &writes {
|
||
match write_one_sidecar(&cache, w) {
|
||
Ok(Outcome::Uploaded) => report.written += 1,
|
||
Ok(Outcome::Queued) => report.queued += 1,
|
||
Err(e) => {
|
||
log::debug!("sidecar for {}: {e}", w.image_path);
|
||
report.last_error = Some(e);
|
||
report.failed += 1;
|
||
}
|
||
}
|
||
}
|
||
report
|
||
};
|
||
|
||
let report = match rt {
|
||
None => queue_all(),
|
||
Some(rt) => rt.block_on(async {
|
||
match NextcloudBackend::new(&creds, &user_id) {
|
||
Ok(b) => {
|
||
let mut report = SidecarReport::default();
|
||
for w in &writes {
|
||
match write_one_sidecar_online(&b, &cache, w).await {
|
||
Ok(Outcome::Uploaded) => report.written += 1,
|
||
Ok(Outcome::Queued) => report.queued += 1,
|
||
Err(e) => {
|
||
log::debug!("sidecar for {}: {e}", w.image_path);
|
||
report.last_error = Some(e);
|
||
report.failed += 1;
|
||
}
|
||
}
|
||
}
|
||
report
|
||
}
|
||
// No backend: the edits are still recorded locally and
|
||
// will go up with the next drain.
|
||
Err(e) => {
|
||
log::debug!("no backend for sidecar upload ({e}); queueing");
|
||
queue_all()
|
||
}
|
||
}
|
||
}),
|
||
};
|
||
|
||
let _ = tx.send(SidecarMessage::Finished {
|
||
written: report.written,
|
||
queued: report.queued,
|
||
failed: report.failed,
|
||
last_error: report.last_error,
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
/// What one write ended up doing.
|
||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||
enum Outcome {
|
||
/// Recorded locally and accepted by the server.
|
||
Uploaded,
|
||
/// Recorded locally, still in the outbox.
|
||
Queued,
|
||
}
|
||
|
||
/// Running totals for a batch, so the loop bodies stay readable.
|
||
#[derive(Debug, Default)]
|
||
struct SidecarReport {
|
||
written: usize,
|
||
queued: usize,
|
||
failed: usize,
|
||
last_error: Option<String>,
|
||
}
|
||
|
||
/// The outcome of a batch of sidecar writes.
|
||
#[derive(Debug)]
|
||
pub enum SidecarMessage {
|
||
Finished {
|
||
written: usize,
|
||
/// Recorded locally but not yet on the server — offline, or an upload
|
||
/// that failed. These are retried by [`spawn_outbox_drain`], so this
|
||
/// is a count of *deferred* work rather than of losses.
|
||
queued: usize,
|
||
failed: usize,
|
||
/// Reported once rather than per file: a network that is down fails
|
||
/// every write with the same message, and forty identical lines in the
|
||
/// status bar say nothing forty times.
|
||
last_error: Option<String>,
|
||
},
|
||
}
|
||
|
||
/// Apply an amendment to a document, returning the new one.
|
||
///
|
||
/// Split out from both write paths so that online and offline produce
|
||
/// *identical* documents: the only thing that differs between them is which
|
||
/// base was read and whether an upload follows. A second copy of this for the
|
||
/// offline case is how the two would come to disagree about what a paste means.
|
||
fn amend(base: dr_pipeline::Sidecar, w: &SidecarWrite) -> dr_pipeline::Sidecar {
|
||
let mut sidecar = base;
|
||
|
||
// Amend the version this write belongs to, creating it if the file did
|
||
// not have one. The uuid comes from the catalog, so the same photograph
|
||
// keeps one identity across devices (FR-NC-8).
|
||
let mut version = sidecar
|
||
.versions
|
||
.get(&w.version_uuid)
|
||
.cloned()
|
||
.unwrap_or_else(|| dr_pipeline::sidecar::Version {
|
||
uuid: w.version_uuid.clone(),
|
||
name: "Default".to_string(),
|
||
is_default: true,
|
||
revision: 0,
|
||
..Default::default()
|
||
});
|
||
|
||
// Only what the amendment names. Everything else in the version — the
|
||
// rating a settings write must not touch, the crop an adjustments-only
|
||
// paste must spare, the unknown keys of an operation this build lacks —
|
||
// survives because it was read from the file and is written back.
|
||
match &w.amendment {
|
||
Amendment::Judgement { rating, flag } => {
|
||
version.rating = *rating;
|
||
version.flag = *flag;
|
||
}
|
||
Amendment::Settings { preset, scope } => {
|
||
preset.amend(&mut version.params, *scope);
|
||
}
|
||
}
|
||
|
||
// A judgement is an edit as far as the merge is concerned, and so is a
|
||
// paste: without the bump, a device that touched the same frame earlier
|
||
// would win on revision and this write would be discarded at the next sync
|
||
// (FR-NC-9).
|
||
version.revision = version.revision.saturating_add(1);
|
||
version.modified = now_secs();
|
||
|
||
sidecar.put(version);
|
||
sidecar
|
||
}
|
||
|
||
/// TRACES: FR-CAT-9
|
||
/// Record an amendment with no server to send it to.
|
||
///
|
||
/// The base is whatever the cache holds, which is either what the server last
|
||
/// had or what earlier offline writes have already built on top of it. Either
|
||
/// way the result is queued, and the drain reconciles it with the server's own
|
||
/// copy when the connection returns — that reconciliation is a *merge*
|
||
/// (FR-NC-9), not an overwrite, so building on a possibly-stale base here does
|
||
/// not cost another device's work.
|
||
fn write_one_sidecar(cache: &SidecarCache, w: &SidecarWrite) -> Result<Outcome, String> {
|
||
let path = sidecar_path(&w.image_path);
|
||
let base = cache.load(&path).unwrap_or_default();
|
||
cache.store(&path, &amend(base, w), true)?;
|
||
Ok(Outcome::Queued)
|
||
}
|
||
|
||
/// TRACES: FR-CAT-8 | FR-CAT-9
|
||
/// Read-modify-write one sidecar, with a server to read from and send to.
|
||
async fn write_one_sidecar_online(
|
||
backend: &NextcloudBackend,
|
||
cache: &SidecarCache,
|
||
w: &SidecarWrite,
|
||
) -> Result<Outcome, String> {
|
||
let path_str = sidecar_path(&w.image_path);
|
||
let path = RemotePath::new(path_str.clone());
|
||
let id = RemoteId::Path(path.clone());
|
||
|
||
// An existing sidecar may hold an edit. Absent is the normal case on a
|
||
// library that has never been edited, and is not an error.
|
||
let existing = backend.get(&id, None).await.ok();
|
||
|
||
// A corrupt sidecar is *not* overwritten: that would destroy an edit this
|
||
// build merely failed to understand. Refused before anything is written,
|
||
// locally or remotely, so the cache cannot end up holding a document that
|
||
// silently discarded the file's real contents.
|
||
if let Some(bytes) = existing.as_deref() {
|
||
if !bytes.is_empty() {
|
||
let text = String::from_utf8_lossy(bytes);
|
||
if dr_pipeline::Sidecar::parse(&text).is_err() {
|
||
return Err(format!("sidecar at {path_str} is unreadable"));
|
||
}
|
||
}
|
||
}
|
||
|
||
let base = existing
|
||
.as_deref()
|
||
.map(|bytes| String::from_utf8_lossy(bytes).into_owned())
|
||
.and_then(|text| dr_pipeline::Sidecar::parse(&text).ok())
|
||
// No sidecar on the server. The cache may still hold queued offline
|
||
// work for this image, and taking `default()` here would drop it.
|
||
.or_else(|| cache.load(&path_str))
|
||
.unwrap_or_default();
|
||
|
||
let sidecar = amend(base, w);
|
||
|
||
// Locally first: this is the commit point, and an upload that fails after
|
||
// it leaves the edit queued rather than lost.
|
||
cache.store(&path_str, &sidecar, true)?;
|
||
|
||
backend
|
||
.put(&path, sidecar.to_text().into_bytes(), None)
|
||
.await
|
||
.map_err(|e| e.to_string())?;
|
||
|
||
// Accepted by the server, so it leaves the outbox. The document stays
|
||
// cached, which is what lets the next offline open still show the edit.
|
||
cache.store(&path_str, &sidecar, false)?;
|
||
Ok(Outcome::Uploaded)
|
||
}
|
||
|
||
/// TRACES: FR-CAT-9 | FR-NC-9 | FR-NC-10
|
||
/// Upload everything the outbox is still holding.
|
||
///
|
||
/// # Why this merges rather than uploads
|
||
///
|
||
/// A queued edit was built on whatever this device last saw. While it sat in
|
||
/// the outbox another device may have edited the same photograph, and simply
|
||
/// PUTting the local document would discard that work — the precise failure
|
||
/// FR-NC-9's node-level merge exists to prevent. So each entry is reconciled
|
||
/// against the server's current copy before it goes up, and disjoint edits
|
||
/// (a crop made here, an exposure change made there) both survive.
|
||
///
|
||
/// # Why an entry stays queued on failure
|
||
///
|
||
/// The marker is cleared only after the server has taken the bytes. A drain
|
||
/// interrupted halfway leaves the rest of the outbox exactly as it was, so
|
||
/// nothing depends on this running to completion.
|
||
pub fn spawn_outbox_drain(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
cache_dir: PathBuf,
|
||
) -> Receiver<SidecarMessage> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
let cache = SidecarCache::open(cache_dir);
|
||
let queued = cache.pending();
|
||
if queued.is_empty() {
|
||
let _ = tx.send(SidecarMessage::Finished {
|
||
written: 0,
|
||
queued: 0,
|
||
failed: 0,
|
||
last_error: None,
|
||
});
|
||
return;
|
||
}
|
||
log::info!("draining {} queued sidecar(s)", queued.len());
|
||
|
||
let rt = match crate::net_runtime::build() {
|
||
Ok(rt) => rt,
|
||
Err(e) => {
|
||
let _ = tx.send(SidecarMessage::Finished {
|
||
written: 0,
|
||
queued: queued.len(),
|
||
failed: 0,
|
||
last_error: Some(e.to_string()),
|
||
});
|
||
return;
|
||
}
|
||
};
|
||
|
||
rt.block_on(async {
|
||
let backend = match NextcloudBackend::new(&creds, &user_id) {
|
||
Ok(b) => b,
|
||
Err(e) => {
|
||
let _ = tx.send(SidecarMessage::Finished {
|
||
written: 0,
|
||
queued: queued.len(),
|
||
failed: 0,
|
||
last_error: Some(e.to_string()),
|
||
});
|
||
return;
|
||
}
|
||
};
|
||
|
||
let mut report = SidecarReport::default();
|
||
for path_str in &queued {
|
||
match drain_one(&backend, &cache, path_str).await {
|
||
Ok(()) => report.written += 1,
|
||
Err(e) => {
|
||
log::debug!("draining {path_str}: {e}");
|
||
report.last_error = Some(e);
|
||
report.failed += 1;
|
||
// Still queued — the marker was never cleared.
|
||
report.queued += 1;
|
||
}
|
||
}
|
||
}
|
||
|
||
let _ = tx.send(SidecarMessage::Finished {
|
||
written: report.written,
|
||
queued: report.queued,
|
||
failed: report.failed,
|
||
last_error: report.last_error,
|
||
});
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
/// Reconcile one queued sidecar with the server and upload it.
|
||
async fn drain_one(
|
||
backend: &NextcloudBackend,
|
||
cache: &SidecarCache,
|
||
path_str: &str,
|
||
) -> Result<(), String> {
|
||
let Some(mut local) = cache.load(path_str) else {
|
||
// The document went while the drain was running. Nothing to send.
|
||
return Ok(());
|
||
};
|
||
|
||
let path = RemotePath::new(path_str.to_string());
|
||
let id = RemoteId::Path(path.clone());
|
||
let remote = backend.get(&id, None).await.ok();
|
||
|
||
if let Some(bytes) = remote.as_deref() {
|
||
if !bytes.is_empty() {
|
||
let text = String::from_utf8_lossy(bytes);
|
||
match dr_pipeline::Sidecar::parse(&text) {
|
||
Ok(remote) => merge_into(&mut local, &remote),
|
||
// Unreadable on the server. Uploading over it would destroy an
|
||
// edit this build failed to understand, so the entry stays
|
||
// queued rather than being resolved destructively.
|
||
Err(e) => return Err(format!("remote sidecar is unreadable ({e})")),
|
||
}
|
||
}
|
||
}
|
||
|
||
backend
|
||
.put(&path, local.to_text().into_bytes(), None)
|
||
.await
|
||
.map_err(|e| e.to_string())?;
|
||
|
||
cache.store(path_str, &local, false)
|
||
}
|
||
|
||
/// TRACES: FR-NC-9
|
||
/// Merge the server's copy into ours, version by version.
|
||
///
|
||
/// No common ancestor is available — the outbox stores the result, not the
|
||
/// base it was built from — so the merge runs with `None`, which treats every
|
||
/// key either side holds as changed. Disjoint keys therefore still both
|
||
/// survive, and a key both sides set resolves by revision exactly as it would
|
||
/// with a base. What is lost without one is the ability to see a *deletion*:
|
||
/// a parameter reset to default on the other device reads as absent rather
|
||
/// than as removed, so our value stands. That is the same direction of caution
|
||
/// the judgement merge takes — an edit is preserved rather than erased.
|
||
fn merge_into(local: &mut dr_pipeline::Sidecar, remote: &dr_pipeline::Sidecar) {
|
||
for (uuid, their_version) in &remote.versions {
|
||
match local.versions.get(uuid).cloned() {
|
||
Some(mut ours) => {
|
||
ours.merge(their_version, None);
|
||
local.put(ours);
|
||
}
|
||
// A version only the server has — another device's virtual copy
|
||
// (FR-CAT-12). Keeping it is what stops one device's upload from
|
||
// deleting another's work.
|
||
None => local.put(their_version.clone()),
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Where the catalog for an account lives.
|
||
///
|
||
/// Keyed by server and user so two accounts do not share an index. Under the
|
||
/// XDG data directory, not cache: the catalog is rebuildable but rebuilding it
|
||
/// costs a full rescan, so it is not something to discard on a cache sweep.
|
||
pub fn catalog_path(server: &str, user_id: &str) -> PathBuf {
|
||
let slug: String = server
|
||
.trim_start_matches("https://")
|
||
.trim_start_matches("http://")
|
||
.chars()
|
||
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
|
||
.collect();
|
||
|
||
let base = std::env::var_os("XDG_DATA_HOME")
|
||
.map(PathBuf::from)
|
||
.or_else(|| std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".local/share")))
|
||
.unwrap_or_else(std::env::temp_dir);
|
||
|
||
base.join("darkroom")
|
||
.join(format!("{slug}-{user_id}"))
|
||
.join("catalog.sqlite")
|
||
}
|
||
|
||
/// 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(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
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, creds, user_id, root, filter, catalog_path, started) {
|
||
let _ = tx.send(ScanMessage::Failed {
|
||
message: e.message,
|
||
offline: e.offline,
|
||
});
|
||
}
|
||
});
|
||
|
||
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.
|
||
struct ScanFailure {
|
||
message: String,
|
||
offline: 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,
|
||
}
|
||
}
|
||
}
|
||
|
||
impl From<dr_sync::RemoteError> for ScanFailure {
|
||
fn from(e: dr_sync::RemoteError) -> Self {
|
||
Self {
|
||
offline: e.indicates_offline(),
|
||
message: e.to_string(),
|
||
}
|
||
}
|
||
}
|
||
|
||
fn run_scan(
|
||
tx: &Sender<ScanMessage>,
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
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 {
|
||
let backend = NextcloudBackend::new(&creds, &user_id).map_err(ScanFailure::local)?;
|
||
|
||
// 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 result = 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?;
|
||
|
||
persist(&catalog, &root, &result).map_err(ScanFailure::local)?;
|
||
|
||
// 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,
|
||
});
|
||
Ok(())
|
||
})
|
||
}
|
||
|
||
/// 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.
|
||
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.
|
||
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()
|
||
});
|
||
|
||
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",
|
||
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(())
|
||
}
|
||
|
||
/// What the grid wants a thumbnail for.
|
||
///
|
||
/// Carries the `oc:fileid` as well as the path, because that is what the
|
||
/// shared store keys on — stable across a server-side move, and the same id
|
||
/// every other client sees (FR-NC-5).
|
||
#[derive(Debug, Clone)]
|
||
pub struct ThumbnailRequest {
|
||
pub row: usize,
|
||
pub path: String,
|
||
/// `None` where the scan found no stable id; such an image is fetched but
|
||
/// not stored, since there is no durable key to store it under.
|
||
pub file_id: Option<u64>,
|
||
/// File length, needed to reject a preview range that points past the end
|
||
/// of the file (NFR-SEC-1).
|
||
pub size: u64,
|
||
/// Catalog row, so EXIF read from the header can be written back.
|
||
pub image_id: i64,
|
||
/// Which resolution this cell needs, from how large it is drawn. A zoomed
|
||
/// grid asks for the large class; a wall of small cells does not.
|
||
///
|
||
/// Named apart from `size`, which is the file's length in bytes — the two
|
||
/// are unrelated and confusing them would fetch the wrong thing.
|
||
pub thumb_size: dr_thumbs::ThumbSize,
|
||
/// Whether this image still needs its EXIF read. Where false the header is
|
||
/// still fetched — the preview needs it — but nothing is parsed or written.
|
||
pub needs_metadata: bool,
|
||
}
|
||
|
||
/// Why a full fetch failed, keeping the one bit the UI cannot re-derive.
|
||
///
|
||
/// The same reasoning as [`ScanFailure`]: the typed error cannot cross the
|
||
/// channel, and "offline" versus "refused" decides whether develop shows
|
||
/// "you are offline — this image is not stored locally" or a real error.
|
||
#[derive(Debug)]
|
||
pub struct FetchFailure {
|
||
pub message: String,
|
||
pub offline: bool,
|
||
}
|
||
|
||
impl FetchFailure {
|
||
fn local(message: impl std::fmt::Display) -> Self {
|
||
Self {
|
||
message: message.to_string(),
|
||
offline: false,
|
||
}
|
||
}
|
||
}
|
||
|
||
impl From<dr_sync::RemoteError> for FetchFailure {
|
||
fn from(e: dr_sync::RemoteError) -> Self {
|
||
Self {
|
||
offline: e.indicates_offline(),
|
||
message: e.to_string(),
|
||
}
|
||
}
|
||
}
|
||
|
||
impl std::fmt::Display for FetchFailure {
|
||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||
f.write_str(&self.message)
|
||
}
|
||
}
|
||
|
||
/// TRACES: FR-NC-6a
|
||
/// Progress from the pin worker.
|
||
#[derive(Debug)]
|
||
pub enum PinMessage {
|
||
/// How many originals the pin still needs. Sent once, before any transfer.
|
||
Planned {
|
||
total: usize,
|
||
},
|
||
/// One original landed.
|
||
Stored {
|
||
done: usize,
|
||
},
|
||
/// The pin is fully downloaded.
|
||
Done {
|
||
stored: usize,
|
||
bytes: u64,
|
||
},
|
||
Failed {
|
||
message: String,
|
||
offline: bool,
|
||
},
|
||
}
|
||
|
||
/// TRACES: FR-NC-6a
|
||
/// Download every original a pin has asked for.
|
||
///
|
||
/// Whole files, deliberately: a pin exists so the photographs can be *edited*
|
||
/// away from the server, and develop needs every photosite. This is the one
|
||
/// place in the app that fetches originals in bulk, which is why FR-NC-6
|
||
/// makes it opt-in rather than something sync does on its own.
|
||
///
|
||
/// Sequential rather than parallel. The lanes that make the thumbnail sweep
|
||
/// fast are wrong here: these are tens of megabytes each, so concurrency buys
|
||
/// little against a single connection's bandwidth and costs a great deal of
|
||
/// memory — and it is the same contention that produced 423 Locked in the
|
||
/// sweep.
|
||
pub fn spawn_pin_fetch(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
catalog_path: PathBuf,
|
||
cache_dir: PathBuf,
|
||
budget: dr_catalog::Budget,
|
||
) -> Receiver<PinMessage> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
let (store, catalog) = match (
|
||
dr_catalog::Cache::open(&cache_dir, budget),
|
||
Catalog::open(&catalog_path),
|
||
) {
|
||
(Ok(s), Ok(c)) => (s, c),
|
||
(Err(e), _) | (_, Err(e)) => {
|
||
let _ = tx.send(PinMessage::Failed {
|
||
message: e.to_string(),
|
||
offline: false,
|
||
});
|
||
return;
|
||
}
|
||
};
|
||
|
||
let pending = match store.pending_pins(catalog.connection()) {
|
||
Ok(p) => p,
|
||
Err(e) => {
|
||
let _ = tx.send(PinMessage::Failed {
|
||
message: e.to_string(),
|
||
offline: false,
|
||
});
|
||
return;
|
||
}
|
||
};
|
||
|
||
if tx
|
||
.send(PinMessage::Planned {
|
||
total: pending.len(),
|
||
})
|
||
.is_err()
|
||
{
|
||
return;
|
||
}
|
||
if pending.is_empty() {
|
||
let _ = tx.send(PinMessage::Done {
|
||
stored: 0,
|
||
bytes: 0,
|
||
});
|
||
return;
|
||
}
|
||
|
||
let rt = match crate::net_runtime::build() {
|
||
Ok(rt) => rt,
|
||
Err(e) => {
|
||
let _ = tx.send(PinMessage::Failed {
|
||
message: e.to_string(),
|
||
offline: false,
|
||
});
|
||
return;
|
||
}
|
||
};
|
||
|
||
rt.block_on(async {
|
||
let backend = match NextcloudBackend::new(&creds, &user_id) {
|
||
Ok(b) => b,
|
||
Err(e) => {
|
||
let _ = tx.send(PinMessage::Failed {
|
||
message: e.to_string(),
|
||
offline: false,
|
||
});
|
||
return;
|
||
}
|
||
};
|
||
|
||
let mut stored = 0usize;
|
||
let mut bytes_total = 0u64;
|
||
|
||
for image in pending {
|
||
let Some(source_ref) = source_ref_of(&catalog, image) else {
|
||
// Catalogued and then removed while the pin was pending.
|
||
continue;
|
||
};
|
||
|
||
let id = RemoteId::Path(RemotePath::new(&source_ref));
|
||
match backend.get(&id, None).await {
|
||
Ok(bytes) => {
|
||
// `pinned: true` — this is the population the budget
|
||
// must never evict, which is the entire promise the
|
||
// user made when they pinned the collection.
|
||
if let Err(e) = store.store(
|
||
catalog.connection(),
|
||
image,
|
||
&source_ref,
|
||
&bytes,
|
||
true,
|
||
now_secs(),
|
||
) {
|
||
log::warn!("storing pinned {source_ref}: {e}");
|
||
continue;
|
||
}
|
||
stored += 1;
|
||
bytes_total += bytes.len() as u64;
|
||
if tx.send(PinMessage::Stored { done: stored }).is_err() {
|
||
return;
|
||
}
|
||
}
|
||
Err(e) if e.indicates_offline() => {
|
||
// Stop rather than failing each remaining file against
|
||
// a dead connection. What was downloaded stays
|
||
// downloaded, and `pending_pins` resumes from there.
|
||
let _ = tx.send(PinMessage::Failed {
|
||
message: e.to_string(),
|
||
offline: true,
|
||
});
|
||
return;
|
||
}
|
||
Err(e) => {
|
||
// One unreadable file must not abandon the whole pin.
|
||
log::warn!("pinning {source_ref}: {e}");
|
||
}
|
||
}
|
||
}
|
||
|
||
let _ = tx.send(PinMessage::Done {
|
||
stored,
|
||
bytes: bytes_total,
|
||
});
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
/// The remote path for a catalogued image.
|
||
fn source_ref_of(catalog: &Catalog, image: dr_types::ImageId) -> Option<String> {
|
||
catalog
|
||
.connection()
|
||
.query_row(
|
||
"SELECT source_ref FROM images WHERE id = ?1",
|
||
rusqlite::params![image.0 as i64],
|
||
|r| r.get(0),
|
||
)
|
||
.ok()
|
||
}
|
||
|
||
/// TRACES: FR-NC-6a | FR-CAT-9
|
||
/// Where a cached original is kept and how much may be kept.
|
||
///
|
||
/// Passed in rather than derived here so the caller owns the policy: the
|
||
/// budget is a user setting, and this function is on a worker thread with no
|
||
/// access to one.
|
||
pub struct CacheContext {
|
||
pub dir: PathBuf,
|
||
pub catalog_path: PathBuf,
|
||
pub image: dr_types::ImageId,
|
||
pub budget: dr_catalog::Budget,
|
||
/// Whether a downloaded original is kept.
|
||
///
|
||
/// Only the write. A cache is always *read*, because bytes already on disk
|
||
/// cost nothing to use and declining them would re-download an image that
|
||
/// is present — including every pinned one, which would leave a pinned
|
||
/// collection unopenable offline the moment this was switched off.
|
||
pub store: bool,
|
||
}
|
||
|
||
/// Fetch one file in full, for opening it in develop.
|
||
///
|
||
/// Deliberately *not* the preview path. Browsing fetches a range and decodes
|
||
/// an embedded JPEG (FR-NC-3); develop needs every byte, because demosaic
|
||
/// needs every photosite. On a RAW file that is tens of megabytes, which is
|
||
/// why this is a click-triggered download and not something the grid does.
|
||
///
|
||
/// # Read-through
|
||
///
|
||
/// With a `cache`, this checks disk before the network and stores what it
|
||
/// downloads. That is what makes opening the same photograph twice cost one
|
||
/// transfer, and what leaves a working session's images openable offline
|
||
/// without anyone having pinned anything.
|
||
///
|
||
/// A cache miss is not an error and a cache failure is not fatal: both fall
|
||
/// through to the network, which is exactly the behaviour that existed before
|
||
/// the cache did.
|
||
///
|
||
/// Returns the bytes on a channel rather than blocking: the download runs on
|
||
/// its own thread and the UI stays live, exactly as thumbnail fetching does.
|
||
/// TRACES: FR-CAT-8 | FR-DEV-6
|
||
/// Fetch and parse the sidecar beside one image.
|
||
///
|
||
/// # Why the edit is read from the file rather than the catalog
|
||
///
|
||
/// The catalog carries a `graph_hash` and no parameters, and it is
|
||
/// *disposable* (ARCH §6.12) — a rebuild would silently return every
|
||
/// photograph to neutral. The sidecar is the authoritative store, so it is
|
||
/// what an open reads, and that is also what makes an edit pasted on the
|
||
/// desktop appear when the same frame is opened on the phone.
|
||
///
|
||
/// # Why absence and failure are the same answer here
|
||
///
|
||
/// `None` means "open this image at its defaults", which is right for a
|
||
/// photograph that has never been edited — the overwhelmingly common case on a
|
||
/// fresh library — and equally right when the network is down. The alternative,
|
||
/// refusing to open the image because its sidecar could not be read, would make
|
||
/// an unreachable server also mean an unviewable library.
|
||
///
|
||
/// The one case that is *not* harmless is a sidecar that exists but does not
|
||
/// parse. That still opens at defaults, but the write path
|
||
/// ([`write_one_sidecar`]) independently refuses to overwrite a file it could
|
||
/// not read, so an edit this build failed to understand is never destroyed by
|
||
/// having been opened.
|
||
pub fn spawn_sidecar_fetch(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
image_path: String,
|
||
cache_dir: PathBuf,
|
||
offline: bool,
|
||
) -> Receiver<Option<dr_pipeline::Sidecar>> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
let cache = SidecarCache::open(cache_dir);
|
||
let path_str = sidecar_path(&image_path);
|
||
|
||
// TRACES: FR-CAT-9 | FR-NC-10
|
||
// The cache wins outright when it is holding work the server has not
|
||
// seen. Fetching in that state would answer with a document *older*
|
||
// than the edit sitting in the outbox, and opening the photograph
|
||
// would silently show it without the change the user just made —
|
||
// which the next save would then write back over the top of.
|
||
if cache.is_pending(&path_str) {
|
||
log::debug!("{path_str} has queued local edits; opening from the cache");
|
||
let _ = tx.send(cache.load(&path_str));
|
||
return;
|
||
}
|
||
|
||
// Offline there is nothing to ask, and the cache is the whole answer.
|
||
let rt = if offline {
|
||
None
|
||
} else {
|
||
match crate::net_runtime::build() {
|
||
Ok(e) => Some(e),
|
||
Err(e) => {
|
||
log::debug!("sidecar fetch runtime: {e}");
|
||
None
|
||
}
|
||
}
|
||
};
|
||
let Some(rt) = rt else {
|
||
let _ = tx.send(cache.load(&path_str));
|
||
return;
|
||
};
|
||
|
||
rt.block_on(async {
|
||
let backend = match NextcloudBackend::new(&creds, &user_id) {
|
||
Ok(b) => b,
|
||
Err(e) => {
|
||
log::debug!("sidecar fetch backend: {e}");
|
||
let _ = tx.send(cache.load(&path_str));
|
||
return;
|
||
}
|
||
};
|
||
|
||
let path = RemotePath::new(path_str.clone());
|
||
let id = RemoteId::Path(path.clone());
|
||
|
||
// A 404 is the normal case on a library that has never been
|
||
// edited, so this is `ok()` rather than an error path.
|
||
let Ok(bytes) = backend.get(&id, None).await else {
|
||
// Unreachable, or no such file. The cache cannot tell those
|
||
// apart and does not need to: either way it holds the best
|
||
// answer this device has.
|
||
let _ = tx.send(cache.load(&path_str));
|
||
return;
|
||
};
|
||
|
||
let text = String::from_utf8_lossy(&bytes).into_owned();
|
||
let parsed = match dr_pipeline::Sidecar::parse(&text) {
|
||
Ok(s) => Some(s),
|
||
Err(e) => {
|
||
log::warn!("sidecar at {} is unreadable ({e})", path.as_str());
|
||
None
|
||
}
|
||
};
|
||
|
||
// Populate the cache from what the server said, so the *next*
|
||
// open of this photograph works with no connection. Clean rather
|
||
// than pending: this content came from the server, so there is
|
||
// nothing to send back.
|
||
if let Some(sidecar) = parsed.as_ref() {
|
||
if let Err(e) = cache.store(&path_str, sidecar, false) {
|
||
log::debug!("caching {path_str}: {e}");
|
||
}
|
||
}
|
||
|
||
let _ = tx.send(parsed);
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
pub fn spawn_full_fetch(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
path: String,
|
||
cache: Option<CacheContext>,
|
||
) -> Receiver<Result<Vec<u8>, FetchFailure>> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
// Opened on this thread: `rusqlite::Connection` is not `Send`, and the
|
||
// UI thread's handle cannot be borrowed across the spawn.
|
||
let cached = cache.as_ref().and_then(|c| {
|
||
let store = dr_catalog::Cache::open(&c.dir, c.budget).ok()?;
|
||
let conn = Catalog::open(&c.catalog_path).ok()?;
|
||
Some((store, conn))
|
||
});
|
||
|
||
if let (Some(c), Some((store, conn))) = (cache.as_ref(), cached.as_ref()) {
|
||
match store.load(conn.connection(), c.image, now_secs()) {
|
||
Ok(Some(bytes)) => {
|
||
log::info!("{path}: {} bytes from the local cache", bytes.len());
|
||
let _ = tx.send(Ok(bytes));
|
||
return;
|
||
}
|
||
Ok(None) => {}
|
||
// A cache that cannot be read is a cache miss, not a failure
|
||
// to open the photograph.
|
||
Err(e) => log::debug!("cache lookup for {path}: {e}"),
|
||
}
|
||
}
|
||
|
||
let rt = match crate::net_runtime::build() {
|
||
Ok(rt) => rt,
|
||
Err(e) => {
|
||
let _ = tx.send(Err(FetchFailure::local(e)));
|
||
return;
|
||
}
|
||
};
|
||
|
||
rt.block_on(async {
|
||
let backend = match NextcloudBackend::new(&creds, &user_id) {
|
||
Ok(b) => b,
|
||
Err(e) => {
|
||
let _ = tx.send(Err(FetchFailure::local(e)));
|
||
return;
|
||
}
|
||
};
|
||
|
||
let id = RemoteId::Path(RemotePath::new(&path));
|
||
let got = backend.get(&id, None).await.map_err(FetchFailure::from);
|
||
|
||
// Store before sending, so the bytes are on disk by the time the
|
||
// image is on screen. Doing it after would leave a window where
|
||
// closing the app immediately lost the download.
|
||
if let (Ok(bytes), Some(c), Some((store, conn))) =
|
||
(&got, cache.as_ref().filter(|c| c.store), cached.as_ref())
|
||
{
|
||
// `pinned: false` — this is the passive population. A pin is
|
||
// something the user asks for explicitly; opening an image is
|
||
// not that, and treating it as one would make the pinned set
|
||
// grow silently and never be evicted.
|
||
if let Err(e) =
|
||
store.store(conn.connection(), c.image, &path, bytes, false, now_secs())
|
||
{
|
||
log::debug!("caching {path}: {e}");
|
||
} else if let Err(e) = store.enforce(conn.connection()) {
|
||
log::debug!("enforcing the cache budget: {e}");
|
||
}
|
||
}
|
||
|
||
let _ = tx.send(got);
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
/// Serve thumbnails for a set of rows: store first, network second.
|
||
///
|
||
/// The store is consulted before any request goes out, so a second launch —
|
||
/// or a second device that synced the shards — fills the grid with no transfer
|
||
/// at all. Only genuine misses reach the network.
|
||
pub fn spawn_thumbnails(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
wanted: Vec<ThumbnailRequest>,
|
||
store_dir: PathBuf,
|
||
catalog_path: PathBuf,
|
||
) -> Receiver<ThumbnailMessage> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
let mut store = match ThumbStore::open(&store_dir) {
|
||
Ok(s) => Some(s),
|
||
Err(e) => {
|
||
// A broken store costs speed, never correctness — every
|
||
// thumbnail can still be fetched.
|
||
log::warn!("thumbnail store unavailable, fetching everything: {e}");
|
||
None
|
||
}
|
||
};
|
||
|
||
// Split the batch before delivering anything, so the plan can be
|
||
// reported first and the UI knows the shape of the work up front.
|
||
// Decoding happens here rather than in the split, because a corrupt
|
||
// blob turns a hit into a miss.
|
||
let mut hits = Vec::new();
|
||
let mut to_fetch = Vec::new();
|
||
// Images whose thumbnail is cached but whose date is still unknown.
|
||
//
|
||
// These need a header read even though no pixels are wanted. Without
|
||
// this pass an image is dated *only* on the one visit that produced
|
||
// its thumbnail — so a library browsed once before the EXIF code
|
||
// existed, or synced from another device's shards, stays permanently
|
||
// undated and never appears on the timeline.
|
||
let mut metadata_only = Vec::new();
|
||
|
||
for req in wanted {
|
||
let stored = req
|
||
.file_id
|
||
.zip(store.as_ref())
|
||
.and_then(|(id, s)| s.get(id, req.thumb_size).ok().flatten());
|
||
|
||
match stored.map(|t| dr_thumbs::decode_rgba(&t.bytes)) {
|
||
Some(Ok((width, height, rgba))) => {
|
||
if req.needs_metadata {
|
||
metadata_only.push(req.clone());
|
||
}
|
||
hits.push(ThumbnailReady {
|
||
row: req.row,
|
||
width,
|
||
height,
|
||
rgba,
|
||
from_cache: true,
|
||
});
|
||
}
|
||
// A corrupt stored blob is a miss, not a failure.
|
||
Some(Err(e)) => {
|
||
log::debug!("stored thumbnail unreadable, refetching: {e}");
|
||
to_fetch.push(req);
|
||
}
|
||
None => to_fetch.push(req),
|
||
}
|
||
}
|
||
|
||
log::info!(
|
||
"thumbnails: {} from store, {} to fetch{}",
|
||
hits.len(),
|
||
to_fetch.len(),
|
||
if metadata_only.is_empty() {
|
||
String::new()
|
||
} else {
|
||
format!(" · {} dates to read", metadata_only.len())
|
||
}
|
||
);
|
||
if tx
|
||
.send(ThumbnailMessage::Plan {
|
||
cached: hits.len(),
|
||
fetching: to_fetch.len(),
|
||
dating: metadata_only.len(),
|
||
})
|
||
.is_err()
|
||
{
|
||
return;
|
||
}
|
||
|
||
for hit in hits {
|
||
if tx.send(ThumbnailMessage::Ready(Box::new(hit))).is_err() {
|
||
return;
|
||
}
|
||
}
|
||
|
||
if to_fetch.is_empty() && metadata_only.is_empty() {
|
||
return;
|
||
}
|
||
|
||
let rt = match crate::net_runtime::build() {
|
||
Ok(rt) => rt,
|
||
Err(e) => {
|
||
for req in &to_fetch {
|
||
let _ = tx.send(ThumbnailMessage::Unavailable {
|
||
row: req.row,
|
||
reason: e.to_string(),
|
||
});
|
||
}
|
||
return;
|
||
}
|
||
};
|
||
|
||
rt.block_on(async {
|
||
let backend = match NextcloudBackend::new(&creds, &user_id) {
|
||
Ok(b) => b,
|
||
Err(e) => {
|
||
for req in &to_fetch {
|
||
let _ = tx.send(ThumbnailMessage::Unavailable {
|
||
row: req.row,
|
||
reason: e.to_string(),
|
||
});
|
||
}
|
||
return;
|
||
}
|
||
};
|
||
|
||
// Batched rather than written per image: one transaction per
|
||
// batch instead of 120, and the grid does not need each date the
|
||
// instant it is read.
|
||
let mut found = Vec::new();
|
||
|
||
// Set when the server proves unreachable, which abandons the rest
|
||
// of the batch. The remaining cells would each take a full timeout
|
||
// to reach the same conclusion — on a 120-cell window, minutes of
|
||
// the grid appearing to load against a server that is not there.
|
||
let mut offline = false;
|
||
|
||
for req in to_fetch {
|
||
let msg = fetch_one(&backend, store.as_mut(), &req, &mut found).await;
|
||
offline = matches!(msg, ThumbnailMessage::Offline { .. });
|
||
// A closed channel means the window went away mid-fetch.
|
||
if tx.send(msg).is_err() || offline {
|
||
break;
|
||
}
|
||
}
|
||
|
||
// Dates for images whose pixels were already cached. Header only —
|
||
// no preview range, no decode.
|
||
//
|
||
// Flushed in chunks rather than once at the end: 119 sequential
|
||
// header fetches take tens of seconds, and a single write at the
|
||
// finish loses every one of them if the window closes first. It
|
||
// also lets the timeline appear while the rest are still arriving.
|
||
//
|
||
// Skipped entirely when the connection has already failed: these
|
||
// are network reads too, and there is nothing left to read from.
|
||
const FLUSH_EVERY: usize = 16;
|
||
if !offline {
|
||
log::info!("reading dates for {} image(s)", metadata_only.len());
|
||
for req in metadata_only {
|
||
if tx.send(ThumbnailMessage::DateProgress).is_err() {
|
||
break;
|
||
}
|
||
read_metadata_only(&backend, &req, &mut found).await;
|
||
|
||
if found.len() >= FLUSH_EVERY {
|
||
flush_metadata(&catalog_path, &mut found, &tx);
|
||
}
|
||
}
|
||
}
|
||
|
||
// Always flushed, even when the batch was abandoned: whatever was
|
||
// read before the connection died is still true, and discarding it
|
||
// would mean re-fetching those headers next time.
|
||
flush_metadata(&catalog_path, &mut found, &tx);
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
/// Fetch a preview in two stages: header, then the exact preview range.
|
||
///
|
||
/// This is what FR-NC-3 specifies, and the single-stage version it replaces
|
||
/// was wrong in a way that looked like corruption: fetching a fixed prefix cut
|
||
/// the embedded JPEG partway through, and decoders render a truncated JPEG as
|
||
/// the top fraction of the frame rather than reporting an error.
|
||
async fn fetch_one(
|
||
backend: &NextcloudBackend,
|
||
store: Option<&mut ThumbStore>,
|
||
req: &ThumbnailRequest,
|
||
found_metadata: &mut Vec<MetadataFound>,
|
||
) -> ThumbnailMessage {
|
||
let id = RemoteId::Path(RemotePath::new(&req.path));
|
||
let fail = |reason: String| ThumbnailMessage::Unavailable {
|
||
row: req.row,
|
||
reason,
|
||
};
|
||
// A connection failure is not this image's verdict. Reported as such so
|
||
// the caller can stop the batch rather than marking sixty cells
|
||
// individually unpreviewable over one dropped connection — a state the
|
||
// grid would then keep until something forced a reload.
|
||
let classify = |e: dr_sync::RemoteError| {
|
||
if e.indicates_offline() {
|
||
ThumbnailMessage::Offline {
|
||
reason: e.to_string(),
|
||
}
|
||
} else {
|
||
fail(e.to_string())
|
||
}
|
||
};
|
||
|
||
// Stage one: the header, enough to parse the container's IFDs.
|
||
let header = match backend.get(&id, Some(0..dr_decode::HEADER_BYTES)).await {
|
||
Ok(b) => b,
|
||
Err(e) => return classify(e),
|
||
};
|
||
|
||
// The same bytes carry EXIF. Reading it here is free — the alternative is
|
||
// a second 256 KB fetch per image over the whole library.
|
||
if req.needs_metadata {
|
||
collect_metadata(&header, req, found_metadata);
|
||
}
|
||
|
||
// Read unconditionally, unlike the rest of the EXIF above: `needs_metadata`
|
||
// is false once an image has been catalogued, but a thumbnail can still be
|
||
// regenerated long after that — a cleared cache, a new size — and a
|
||
// thumbnail that came out upright the first time must come out upright
|
||
// every time. This is a header walk, not a decode; see `dr_decode::orientation`.
|
||
let orientation = dr_decode::orientation(&header).unwrap_or_default();
|
||
|
||
// A plain JPEG is its own preview; anything else needs locating.
|
||
let bytes = if header.starts_with(&[0xFF, 0xD8, 0xFF]) {
|
||
match backend.get(&id, None).await {
|
||
Ok(b) => b,
|
||
Err(e) => return classify(e),
|
||
}
|
||
} else {
|
||
let Some(loc) = dr_decode::locate_preview(&header, req.size) else {
|
||
// No locatable preview. Declining beats fetching the whole file:
|
||
// that is the 370 GB path FR-NC-3 exists to avoid.
|
||
return fail("no locatable embedded preview".into());
|
||
};
|
||
if loc.len() > MAX_PREVIEW_BYTES {
|
||
return fail(format!("preview is {} bytes, too large", loc.len()));
|
||
}
|
||
|
||
// Stage two: exactly the preview's bytes.
|
||
match backend.get(&id, Some(loc.range.clone())).await {
|
||
Ok(b) => b,
|
||
Err(e) => return classify(e),
|
||
}
|
||
};
|
||
|
||
// Verify before decoding. A truncated JPEG decodes "successfully" into a
|
||
// partial frame, so without this the broken result reaches the cache and
|
||
// the screen looking like a corrupt file.
|
||
if !dr_decode::is_complete_jpeg(&bytes) {
|
||
return fail("preview bytes are incomplete".into());
|
||
}
|
||
|
||
// Decode on the worker, never the UI thread.
|
||
let mut preview = match dr_decode::decode_jpeg(&bytes) {
|
||
Ok(p) => p,
|
||
Err(e) => return fail(e.to_string()),
|
||
};
|
||
preview.downscale_to(req.thumb_size.edge());
|
||
// Turn it the right way up before it is measured, cached or shown. An
|
||
// embedded preview is written in the sensor's orientation, so without this
|
||
// every frame shot in portrait lies on its side in the grid — and, because
|
||
// the cache is keyed by file and size alone, stays that way.
|
||
//
|
||
// After the downscale, so the permutation moves thumbnail-sized bytes
|
||
// rather than the full preview's.
|
||
preview.apply_orientation(orientation);
|
||
|
||
// Persist for next time, and for every other client that syncs the shard.
|
||
// A store failure is logged and dropped: the pixels are already in hand,
|
||
// and refusing to display them because they could not be cached would be
|
||
// the wrong trade.
|
||
if let (Some(store), Some(file_id)) = (store, req.file_id) {
|
||
match dr_thumbs::encode_rgba(preview.width, preview.height, &preview.rgba) {
|
||
Ok(encoded) => {
|
||
let thumb = dr_thumbs::Thumbnail {
|
||
width: preview.width,
|
||
height: preview.height,
|
||
bytes: encoded,
|
||
};
|
||
if let Err(e) = store.put(file_id, req.thumb_size, &thumb) {
|
||
log::debug!("storing thumbnail {file_id}: {e}");
|
||
}
|
||
}
|
||
Err(e) => log::debug!("encoding thumbnail {file_id}: {e}"),
|
||
}
|
||
}
|
||
|
||
ThumbnailMessage::Ready(Box::new(ThumbnailReady {
|
||
row: req.row,
|
||
width: preview.width,
|
||
height: preview.height,
|
||
rgba: preview.rgba,
|
||
from_cache: false,
|
||
}))
|
||
}
|
||
|
||
/// Parse EXIF out of a header and record it.
|
||
///
|
||
/// Shared by both paths — the thumbnail fetch, which gets the header anyway,
|
||
/// and the header-only pass for images whose pixels were already cached.
|
||
fn collect_metadata(header: &[u8], req: &ThumbnailRequest, out: &mut Vec<MetadataFound>) {
|
||
let Ok(md) = dr_decode::metadata(header) else {
|
||
return;
|
||
};
|
||
out.push(MetadataFound {
|
||
image_id: req.image_id,
|
||
captured_at: md.captured_at,
|
||
captured_offset: md.captured_offset,
|
||
camera: match (&md.make, &md.model) {
|
||
// Bodies repeat the make inside the model ("Canon EOS 6D"), so
|
||
// joining unconditionally yields "Canon Canon EOS 6D".
|
||
(Some(make), Some(model)) if model.starts_with(make.as_str()) => {
|
||
Some(model.trim().to_string())
|
||
}
|
||
(Some(make), Some(model)) => Some(format!("{} {}", make.trim(), model.trim())),
|
||
(None, Some(model)) => Some(model.trim().to_string()),
|
||
_ => None,
|
||
},
|
||
lens: md.lens.map(|l| l.trim().to_string()),
|
||
iso: md.iso,
|
||
});
|
||
}
|
||
|
||
/// Write a batch of dates and tell the UI, draining `found`.
|
||
///
|
||
/// Separate from the loop so the same path serves both the periodic flush and
|
||
/// the final one, and so a write failure is reported once rather than being
|
||
/// silently swallowed by the caller.
|
||
fn flush_metadata(
|
||
catalog_path: &std::path::Path,
|
||
found: &mut Vec<MetadataFound>,
|
||
tx: &Sender<ThumbnailMessage>,
|
||
) {
|
||
if found.is_empty() {
|
||
return;
|
||
}
|
||
match Catalog::open(catalog_path) {
|
||
Ok(cat) => match write_metadata(&cat, found) {
|
||
Ok(n) => {
|
||
log::info!("recorded capture dates for {n} of {} image(s)", found.len());
|
||
// Tell the UI so the timeline can appear. Without this the
|
||
// histogram only shows up on the next window load, which on a
|
||
// fully cached library may be never.
|
||
let _ = tx.send(ThumbnailMessage::DatesRecorded(n));
|
||
}
|
||
Err(e) => log::warn!("writing metadata: {e}"),
|
||
},
|
||
Err(e) => log::warn!("opening catalog to write metadata: {e}"),
|
||
}
|
||
found.clear();
|
||
}
|
||
|
||
/// Read only the date for an image whose thumbnail is already cached.
|
||
///
|
||
/// One 256 KB header request, no preview range and no decode. This is what
|
||
/// gets a library dated when its thumbnails came from the store — including
|
||
/// shards synced from another device, which carry pixels but no metadata.
|
||
/// Read a header for its date.
|
||
///
|
||
/// Returns whether the file was **reached**, which the caller needs and cannot
|
||
/// otherwise tell: a header that carried no EXIF and a fetch that never
|
||
/// happened both leave `found` untouched, and recording the second as "this
|
||
/// image has no date" would let one lock mark it dateless for good.
|
||
async fn read_metadata_only(
|
||
backend: &NextcloudBackend,
|
||
req: &ThumbnailRequest,
|
||
found: &mut Vec<MetadataFound>,
|
||
) -> bool {
|
||
let id = RemoteId::Path(RemotePath::new(&req.path));
|
||
|
||
// Retried, because one failure here is usually a lock rather than a
|
||
// verdict. Nextcloud's file locking answers a plain *read* with 423 under
|
||
// concurrency, and the identical range succeeds moments later — measured
|
||
// against a real server while twelve lanes were running. Without a retry
|
||
// those images sit out the whole pass over a lock that lasted a moment.
|
||
//
|
||
// Bounded and short: a genuinely missing or forbidden file must not cost
|
||
// three round trips before the sweep moves on.
|
||
const ATTEMPTS: usize = 3;
|
||
for attempt in 1..=ATTEMPTS {
|
||
match backend.get(&id, Some(0..dr_decode::HEADER_BYTES)).await {
|
||
Ok(header) => {
|
||
collect_metadata(&header, req, found);
|
||
return true;
|
||
}
|
||
Err(e) if e.is_transient() && attempt < ATTEMPTS => {
|
||
// Backing off at all matters more than the exact interval: the
|
||
// contention that produced the lock is our own lanes, so any
|
||
// pause lets the holder finish.
|
||
tokio::time::sleep(std::time::Duration::from_millis(200 * attempt as u64)).await;
|
||
}
|
||
Err(e) => {
|
||
// Not surfaced: a missing date leaves the image off the
|
||
// timeline rather than breaking anything, and the next sweep
|
||
// retries it regardless.
|
||
log::debug!("reading date for {} ({attempt} attempts): {e}", req.path);
|
||
return false;
|
||
}
|
||
}
|
||
}
|
||
false
|
||
}
|
||
|
||
/// Write capture metadata read during the thumbnail pass.
|
||
///
|
||
/// Promotes each row from `metadata_state = 1` (stat-only) to 2 (full EXIF),
|
||
/// which is what makes it eligible for the timeline. A row whose EXIF was
|
||
/// unreadable stays at 1 rather than being marked done with empty fields, so a
|
||
/// later attempt can retry it.
|
||
///
|
||
/// Returns how many rows were promoted.
|
||
pub fn write_metadata(
|
||
catalog: &Catalog,
|
||
found: &[MetadataFound],
|
||
) -> Result<usize, dr_catalog::CatalogError> {
|
||
let conn = catalog.connection();
|
||
let tx = conn.unchecked_transaction()?;
|
||
let mut promoted = 0;
|
||
|
||
for m in found {
|
||
// Only a real timestamp counts as fully read. Camera and lens without
|
||
// a date leave the image unplaceable on a timeline, which is exactly
|
||
// the state the grid needs to distinguish.
|
||
let state = if m.captured_at.is_some() { 2 } else { 1 };
|
||
|
||
tx.execute(
|
||
"UPDATE images
|
||
SET captured_at = coalesce(?2, captured_at),
|
||
captured_offset = coalesce(?3, captured_offset),
|
||
camera = coalesce(?4, camera),
|
||
lens = coalesce(?5, lens),
|
||
iso = coalesce(?6, iso),
|
||
metadata_state = max(metadata_state, ?7)
|
||
WHERE id = ?1",
|
||
rusqlite::params![
|
||
m.image_id,
|
||
m.captured_at,
|
||
m.captured_offset,
|
||
m.camera,
|
||
m.lens,
|
||
m.iso,
|
||
state,
|
||
],
|
||
)?;
|
||
if state == 2 {
|
||
promoted += 1;
|
||
}
|
||
}
|
||
|
||
tx.commit()?;
|
||
Ok(promoted)
|
||
}
|
||
|
||
/// Progress from the whole-library sweep.
|
||
#[derive(Debug)]
|
||
pub enum SweepMessage {
|
||
/// How many images still need work, counted once at the start.
|
||
Total(usize),
|
||
/// Another chunk finished. Carries cumulative counts.
|
||
Progress {
|
||
done: usize,
|
||
dated: usize,
|
||
},
|
||
Finished {
|
||
dated: usize,
|
||
},
|
||
}
|
||
|
||
/// Await every future concurrently, returning results in order.
|
||
///
|
||
/// A hand-rolled `join_all` rather than a `futures` dependency for one
|
||
/// function. Polling a `Vec` of futures in a loop is exactly what the crate's
|
||
/// version does; the ordering guarantee is what lets the caller pair results
|
||
/// back to their inputs.
|
||
async fn futures_join_all<F>(futures: impl IntoIterator<Item = F>) -> Vec<F::Output>
|
||
where
|
||
F: std::future::Future,
|
||
{
|
||
use std::pin::Pin;
|
||
use std::task::Poll;
|
||
|
||
// Boxed so each future has a stable address while it is polled in place.
|
||
let mut pending: Vec<Option<Pin<Box<F>>>> =
|
||
futures.into_iter().map(|f| Some(Box::pin(f))).collect();
|
||
let mut done: Vec<Option<F::Output>> = (0..pending.len()).map(|_| None).collect();
|
||
|
||
std::future::poll_fn(move |cx| {
|
||
let mut all_ready = true;
|
||
for (slot, out) in pending.iter_mut().zip(done.iter_mut()) {
|
||
let Some(fut) = slot else { continue };
|
||
match fut.as_mut().poll(cx) {
|
||
Poll::Ready(v) => {
|
||
*out = Some(v);
|
||
// Dropped as soon as it completes, so a long-running lane
|
||
// does not hold a finished one's resources.
|
||
*slot = None;
|
||
}
|
||
Poll::Pending => all_ready = false,
|
||
}
|
||
}
|
||
if all_ready {
|
||
Poll::Ready(done.iter_mut().filter_map(Option::take).collect())
|
||
} else {
|
||
Poll::Pending
|
||
}
|
||
})
|
||
.await
|
||
}
|
||
|
||
/// How many images one sweep chunk handles before committing.
|
||
///
|
||
/// Small enough that a kill loses little, large enough that the catalog is not
|
||
/// reopened per image. A multiple of [`SWEEP_LANES`] so every lane gets equal
|
||
/// work and no chunk ends with most lanes idle.
|
||
const SWEEP_CHUNK: usize = 96;
|
||
|
||
/// How many fetches the sweep keeps in flight.
|
||
///
|
||
/// Each is ~0.6 s of round-trip latency and almost no bandwidth — a 256 KB
|
||
/// header — so the sequential version spent essentially all its time waiting.
|
||
/// Twelve lanes turn ~3 hours into ~15 minutes on the reference library.
|
||
///
|
||
/// Deliberately bounded rather than unlimited: the grid's own interactive
|
||
/// fetches share this server, and a sweep that saturated the connection would
|
||
/// make browsing feel broken while it ran.
|
||
///
|
||
/// **Lowered from twelve after measuring.** Twelve produced 423 Locked on a
|
||
/// real server — Nextcloud's file locking answering a plain read under
|
||
/// contention we were creating ourselves. Six keeps most of the speedup
|
||
/// without provoking it; the retry above covers what still slips through.
|
||
const SWEEP_LANES: usize = 6;
|
||
|
||
/// Date **every** image in the library, not just the ones on screen.
|
||
///
|
||
/// Thumbnails are deliberately *not* fetched here. A thumbnail needs the
|
||
/// mutable store, which cannot be shared across the parallel lanes below, and
|
||
/// it costs 1–3 MB against a date's 256 KB. Dating the whole library is what
|
||
/// the timeline needs; thumbnails arrive as cells are actually browsed, which
|
||
/// is the FR-NC-3 posture anyway.
|
||
///
|
||
/// The grid's own fetches cover what is on screen; this covers the rest, so the
|
||
/// timeline describes the whole library rather than the part that happened to
|
||
/// be scrolled past. It is resumable by construction — each pass queries for
|
||
/// what is still missing, so a kill mid-sweep costs only the current chunk.
|
||
///
|
||
/// Runs at the back of the queue by design: it holds no lock the grid needs,
|
||
/// and its chunked commits keep write transactions short.
|
||
pub fn spawn_sweep(
|
||
creds: AppCredentials,
|
||
user_id: String,
|
||
catalog_path: PathBuf,
|
||
) -> Receiver<SweepMessage> {
|
||
let (tx, rx) = std::sync::mpsc::channel();
|
||
|
||
std::thread::spawn(move || {
|
||
let catalog = match Catalog::open(&catalog_path) {
|
||
Ok(c) => c,
|
||
Err(e) => {
|
||
// Silent failure here left the sweep looking like it had run
|
||
// and found nothing: no progress, no error, 17,397 images
|
||
// still unindexed.
|
||
log::warn!(
|
||
"sweep: cannot open catalog at {}: {e}",
|
||
catalog_path.display()
|
||
);
|
||
let _ = tx.send(SweepMessage::Finished { dated: 0 });
|
||
return;
|
||
}
|
||
};
|
||
|
||
let outstanding = count_outstanding(&catalog).unwrap_or(0);
|
||
if outstanding == 0 {
|
||
let _ = tx.send(SweepMessage::Finished { dated: 0 });
|
||
return;
|
||
}
|
||
log::info!("sweep: {outstanding} image(s) need a date or a thumbnail");
|
||
if tx.send(SweepMessage::Total(outstanding)).is_err() {
|
||
return;
|
||
}
|
||
|
||
let rt = match crate::net_runtime::build() {
|
||
Ok(rt) => rt,
|
||
Err(e) => {
|
||
log::warn!("sweep: no runtime: {e}");
|
||
return;
|
||
}
|
||
};
|
||
|
||
rt.block_on(async {
|
||
let Ok(backend) = NextcloudBackend::new(&creds, &user_id) else {
|
||
return;
|
||
};
|
||
|
||
let (mut done, mut dated) = (0usize, 0usize);
|
||
loop {
|
||
// Re-queried each pass rather than held as one long list: the
|
||
// grid is dating images at the same time, and a stale list
|
||
// would refetch what it already covered.
|
||
let chunk = match next_outstanding(&catalog, SWEEP_CHUNK) {
|
||
Ok(c) if c.is_empty() => break,
|
||
Ok(c) => c,
|
||
Err(e) => {
|
||
log::warn!("sweep: {e}");
|
||
break;
|
||
}
|
||
};
|
||
|
||
let chunk_started = std::time::Instant::now();
|
||
log::debug!(
|
||
"sweep: chunk of {} starting at image {}",
|
||
chunk.len(),
|
||
chunk[0].image_id
|
||
);
|
||
|
||
// Twelve lanes over the chunk. Each lane owns a disjoint slice
|
||
// and its own `found` vector, so nothing is shared and no lock
|
||
// is needed; the results are concatenated after the join.
|
||
//
|
||
// The thumbnail store is the exception — it is `&mut` and
|
||
// cannot be shared — so lanes only *read* metadata and any
|
||
// missing thumbnail is left to the interactive path. Dating the
|
||
// library is what the sweep is for; thumbnails arrive as cells
|
||
// are browsed.
|
||
let lanes: Vec<Vec<&ThumbnailRequest>> = (0..SWEEP_LANES)
|
||
.map(|lane| chunk.iter().skip(lane).step_by(SWEEP_LANES).collect())
|
||
.collect();
|
||
|
||
let results = futures_join_all(lanes.into_iter().map(|lane| {
|
||
let backend = &backend;
|
||
async move {
|
||
let mut found = Vec::new();
|
||
let mut reached = Vec::new();
|
||
for req in lane {
|
||
if read_metadata_only(backend, req, &mut found).await {
|
||
reached.push(req.image_id);
|
||
}
|
||
}
|
||
(found, reached)
|
||
}
|
||
}))
|
||
.await;
|
||
|
||
let mut found = Vec::new();
|
||
let mut reached = std::collections::HashSet::new();
|
||
for (lane_found, lane_reached) in results {
|
||
found.extend(lane_found);
|
||
reached.extend(lane_reached);
|
||
}
|
||
done += chunk.len();
|
||
|
||
// An image whose header carried no EXIF at all yields nothing
|
||
// to `found`, so nothing marks it examined and the next sweep
|
||
// fetches it again — for ever. Darktable exports strip
|
||
// metadata by default, and 2,188 of them in the reference
|
||
// library meant 2,188 pointless round trips per run.
|
||
//
|
||
// Recorded as examined with no date: the file was read and
|
||
// genuinely has none, which is a different state from "not
|
||
// looked at yet" and must not be confused with it.
|
||
let answered: std::collections::HashSet<i64> =
|
||
found.iter().map(|m| m.image_id).collect();
|
||
// Only files actually read. One that could not be fetched is
|
||
// left alone so the next pass retries it, rather than being
|
||
// written off over a lock or a dropped connection.
|
||
found.extend(
|
||
chunk
|
||
.iter()
|
||
.filter(|r| {
|
||
reached.contains(&r.image_id) && !answered.contains(&r.image_id)
|
||
})
|
||
.map(|r| MetadataFound {
|
||
image_id: r.image_id,
|
||
captured_at: None,
|
||
captured_offset: None,
|
||
camera: None,
|
||
lens: None,
|
||
iso: None,
|
||
}),
|
||
);
|
||
|
||
dated += found.iter().filter(|m| m.captured_at.is_some()).count();
|
||
let read = answered.len();
|
||
flush_sweep(&catalog, &mut found);
|
||
log::info!(
|
||
"sweep: {done} done, {dated} dated ({read} read in {:.1}s)",
|
||
chunk_started.elapsed().as_secs_f64()
|
||
);
|
||
|
||
if tx.send(SweepMessage::Progress { done, dated }).is_err() {
|
||
return;
|
||
}
|
||
}
|
||
|
||
log::info!("sweep complete: {dated} date(s) recorded over {done} image(s)");
|
||
let _ = tx.send(SweepMessage::Finished { dated });
|
||
});
|
||
});
|
||
|
||
rx
|
||
}
|
||
|
||
/// How many images still lack a date or a thumbnail.
|
||
fn count_outstanding(catalog: &Catalog) -> Result<usize, dr_catalog::CatalogError> {
|
||
let n: i64 = catalog.connection().query_row(
|
||
&format!(
|
||
"SELECT count(*) FROM images
|
||
WHERE metadata_state < 2 AND {VISIBLE_UNALIASED}"
|
||
),
|
||
[],
|
||
|r| r.get(0),
|
||
)?;
|
||
Ok(n as usize)
|
||
}
|
||
|
||
/// The next images needing work.
|
||
///
|
||
/// Ordered by id so the sweep advances deterministically and a resumed run
|
||
/// picks up where it left off rather than revisiting.
|
||
fn next_outstanding(
|
||
catalog: &Catalog,
|
||
limit: usize,
|
||
) -> Result<Vec<ThumbnailRequest>, dr_catalog::CatalogError> {
|
||
let mut stmt = catalog.connection().prepare(&format!(
|
||
"SELECT i.id, i.source_ref, r.file_id, i.file_size
|
||
FROM images i
|
||
LEFT JOIN remote r ON r.image_id = i.id
|
||
WHERE i.metadata_state < 2 AND {VISIBLE}
|
||
ORDER BY i.id
|
||
LIMIT ?1"
|
||
))?;
|
||
let rows = stmt
|
||
.query_map([limit as i64], |r| {
|
||
Ok(ThumbnailRequest {
|
||
// The sweep indexes dates, and reads headers only — the size
|
||
// never reaches a fetch, but it must name something.
|
||
thumb_size: dr_thumbs::ThumbSize::Grid,
|
||
// Row index is meaningless here — the sweep touches no grid
|
||
// cell, so nothing consumes it.
|
||
row: 0,
|
||
image_id: r.get(0)?,
|
||
path: r.get(1)?,
|
||
file_id: r.get::<_, Option<i64>>(2)?.map(|v| v as u64),
|
||
size: r.get::<_, Option<i64>>(3)?.unwrap_or(0) as u64,
|
||
needs_metadata: true,
|
||
})
|
||
})?
|
||
.collect::<Result<Vec<_>, _>>()?;
|
||
Ok(rows)
|
||
}
|
||
|
||
/// Commit a sweep chunk.
|
||
///
|
||
/// An image whose header yielded no date is still marked done, or the sweep
|
||
/// would revisit it forever. `write_metadata` records `metadata_state = 1` for
|
||
/// those, so this promotes them explicitly.
|
||
fn flush_sweep(catalog: &Catalog, found: &mut Vec<MetadataFound>) {
|
||
if found.is_empty() {
|
||
return;
|
||
}
|
||
if let Err(e) = write_metadata(catalog, found) {
|
||
log::warn!("sweep: writing metadata: {e}");
|
||
found.clear();
|
||
return;
|
||
}
|
||
|
||
// Mark the dateless as examined. Without this they stay at state 1 and the
|
||
// sweep loops over them on every pass, never terminating.
|
||
let ids: Vec<i64> = found
|
||
.iter()
|
||
.filter(|m| m.captured_at.is_none())
|
||
.map(|m| m.image_id)
|
||
.collect();
|
||
for id in ids {
|
||
let _ = catalog
|
||
.connection()
|
||
.execute("UPDATE images SET metadata_state = 2 WHERE id = ?1", [id]);
|
||
}
|
||
found.clear();
|
||
}
|
||
|
||
/// Where an account's thumbnail shards live.
|
||
///
|
||
/// Beside the catalog rather than in the cache directory: these sync to the
|
||
/// server and are shared with other clients, so discarding them on a cache
|
||
/// sweep would cost a re-download for everyone.
|
||
pub fn thumbs_dir(server: &str, user_id: &str) -> PathBuf {
|
||
catalog_path(server, user_id)
|
||
.parent()
|
||
.map(|p| p.join("thumbs"))
|
||
.unwrap_or_else(|| std::env::temp_dir().join("darkroom-thumbs"))
|
||
}
|
||
|
||
/// One grid cell's data, read from the catalog.
|
||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||
pub struct LibraryCell {
|
||
pub image_id: i64,
|
||
pub name: String,
|
||
pub remote_path: String,
|
||
/// `oc:fileid`, the key the shared thumbnail store uses. `None` for an
|
||
/// image the scan found without a stable id.
|
||
pub file_id: Option<u64>,
|
||
/// File length, for bounds-checking a located preview range.
|
||
pub size: u64,
|
||
/// 0 = nothing, 1 = stat-only, 2 = full EXIF.
|
||
pub metadata_state: u8,
|
||
/// UTC seconds, once EXIF has been read.
|
||
pub captured_at: Option<i64>,
|
||
}
|
||
|
||
/// Read a window of cells out of the catalog.
|
||
///
|
||
/// Windowed rather than wholesale: a 17k-image library must not become 17k
|
||
/// rows in a Slint model (FR-CAT-4).
|
||
/// The unscoped form, kept as the name the tests and any future caller reach
|
||
/// for. The UI goes through [`read_cells_scoped`], because a collection may be
|
||
/// selected.
|
||
#[cfg(test)]
|
||
pub fn read_cells(
|
||
catalog: &Catalog,
|
||
offset: usize,
|
||
limit: usize,
|
||
) -> Result<Vec<LibraryCell>, dr_catalog::CatalogError> {
|
||
read_cells_scoped(catalog, None, &RatingFilter::default(), offset, limit)
|
||
}
|
||
|
||
/// Read a window of cells, optionally narrowed to one collection.
|
||
///
|
||
/// A collection *set* shows its descendants' images too — a parent whose
|
||
/// children hold everything would otherwise read as empty, which makes nesting
|
||
/// look broken. The id list comes from
|
||
/// [`dr_catalog::collections::descendants`], which is depth-guarded.
|
||
///
|
||
/// Ordering matches the unscoped grid (capture time, then name) rather than
|
||
/// manual position: position is only meaningful inside one collection and this
|
||
/// query also serves sets, where two children's positions are unrelated.
|
||
pub fn read_cells_scoped(
|
||
catalog: &Catalog,
|
||
scope: Option<dr_types::CollectionId>,
|
||
filter: &RatingFilter,
|
||
offset: usize,
|
||
limit: usize,
|
||
) -> Result<Vec<LibraryCell>, dr_catalog::CatalogError> {
|
||
let Some(scope) = scope else {
|
||
return read_cells_all(catalog, filter, offset, limit);
|
||
};
|
||
|
||
let ids = dr_catalog::collections::descendants(catalog.connection(), scope)?;
|
||
// Placeholders are generated from the *count* of ids, never from user text.
|
||
let placeholders = std::iter::repeat_n("?", ids.len())
|
||
.collect::<Vec<_>>()
|
||
.join(",");
|
||
let rated = filter.sql();
|
||
let sql = format!(
|
||
"SELECT i.id, i.source_ref, r.file_id, i.file_size,
|
||
i.metadata_state, i.captured_at
|
||
FROM images i
|
||
LEFT JOIN remote r ON r.image_id = i.id
|
||
WHERE {VISIBLE}{rated}
|
||
AND i.id IN (SELECT image_id FROM collection_members
|
||
WHERE collection_id IN ({placeholders}))
|
||
ORDER BY i.captured_at IS NULL, i.captured_at ASC, i.source_ref ASC
|
||
LIMIT ? OFFSET ?"
|
||
);
|
||
|
||
let mut params: Vec<rusqlite::types::Value> = ids
|
||
.iter()
|
||
.map(|c| rusqlite::types::Value::Integer(c.0 as i64))
|
||
.collect();
|
||
params.push(rusqlite::types::Value::Integer(limit as i64));
|
||
params.push(rusqlite::types::Value::Integer(offset as i64));
|
||
|
||
let mut stmt = catalog.connection().prepare(&sql)?;
|
||
let rows = stmt
|
||
.query_map(rusqlite::params_from_iter(params.iter()), row_to_cell)?
|
||
.collect::<Result<Vec<_>, _>>()?;
|
||
Ok(rows)
|
||
}
|
||
|
||
fn read_cells_all(
|
||
catalog: &Catalog,
|
||
filter: &RatingFilter,
|
||
offset: usize,
|
||
limit: usize,
|
||
) -> Result<Vec<LibraryCell>, dr_catalog::CatalogError> {
|
||
let rated = filter.sql();
|
||
let mut stmt = catalog.connection().prepare(&format!(
|
||
"SELECT i.id, i.source_ref, r.file_id, i.file_size,
|
||
i.metadata_state, i.captured_at
|
||
FROM images i
|
||
LEFT JOIN remote r ON r.image_id = i.id
|
||
WHERE {VISIBLE}{rated}
|
||
ORDER BY i.captured_at IS NULL, i.captured_at ASC, i.source_ref ASC
|
||
LIMIT ?1 OFFSET ?2"
|
||
))?;
|
||
let rows = stmt
|
||
.query_map(rusqlite::params![limit as i64, offset as i64], row_to_cell)?
|
||
.collect::<Result<Vec<_>, _>>()?;
|
||
Ok(rows)
|
||
}
|
||
|
||
/// TRACES: FR-CAT-15
|
||
/// Read a window of *trashed* cells, newest deletion first.
|
||
///
|
||
/// Ordered by when it was trashed rather than by capture time, which is what
|
||
/// every other view sorts by. The question in the trash is "what did I just
|
||
/// delete?", not "when was this taken" — a mistaken delete is corrected within
|
||
/// seconds, and burying it among photographs from the same afternoon would make
|
||
/// the one row the user is looking for the hardest one to find.
|
||
///
|
||
/// The rating filter is deliberately not applied. It narrows a *culling* pass,
|
||
/// and a trash that hid rows because of a filter set elsewhere would look like
|
||
/// it had lost them.
|
||
pub fn read_trashed_cells(
|
||
catalog: &Catalog,
|
||
offset: usize,
|
||
limit: usize,
|
||
) -> Result<Vec<LibraryCell>, dr_catalog::CatalogError> {
|
||
let mut stmt = catalog.connection().prepare(&format!(
|
||
"SELECT i.id, i.source_ref, r.file_id, i.file_size,
|
||
i.metadata_state, i.captured_at
|
||
FROM images i
|
||
LEFT JOIN remote r ON r.image_id = i.id
|
||
WHERE {TRASHED}
|
||
ORDER BY i.trashed_at DESC, i.source_ref ASC
|
||
LIMIT ?1 OFFSET ?2"
|
||
))?;
|
||
let rows = stmt
|
||
.query_map(rusqlite::params![limit as i64, offset as i64], row_to_cell)?
|
||
.collect::<Result<Vec<_>, _>>()?;
|
||
Ok(rows)
|
||
}
|
||
|
||
/// TRACES: FR-CAT-15
|
||
/// How many images the trash view would list.
|
||
///
|
||
/// Counts exactly what [`read_trashed_cells`] lists — same predicate, no filter
|
||
/// — so the scrollbar and the header cannot disagree with the cells.
|
||
pub fn total_trashed(catalog: &Catalog) -> Result<usize, dr_catalog::CatalogError> {
|
||
let n: i64 = catalog.connection().query_row(
|
||
&format!("SELECT count(*) FROM images i WHERE {TRASHED}"),
|
||
[],
|
||
|r| r.get(0),
|
||
)?;
|
||
Ok(n as usize)
|
||
}
|
||
|
||
/// Shared row mapping, so the scoped and unscoped queries cannot drift.
|
||
fn row_to_cell(r: &rusqlite::Row) -> rusqlite::Result<LibraryCell> {
|
||
let path: String = r.get(1)?;
|
||
Ok(LibraryCell {
|
||
image_id: r.get(0)?,
|
||
name: path.rsplit(['/', ':']).next().unwrap_or(&path).to_string(),
|
||
remote_path: path,
|
||
file_id: r.get::<_, Option<i64>>(2)?.map(|v| v as u64),
|
||
size: r.get::<_, Option<i64>>(3)?.unwrap_or(0) as u64,
|
||
metadata_state: r.get::<_, i64>(4)? as u8,
|
||
captured_at: r.get(5)?,
|
||
})
|
||
}
|
||
|
||
/// Total images in the catalog, or in one collection and its descendants.
|
||
///
|
||
/// Counts exactly what [`read_cells_scoped`] would list, filter included. The
|
||
/// two must agree: the header says "412 images" and the grid's scrollbar is
|
||
/// sized from the same number, so a count that ignored the filter would leave
|
||
/// the user scrolling through empty rows.
|
||
pub fn total_images_scoped(
|
||
catalog: &Catalog,
|
||
scope: Option<dr_types::CollectionId>,
|
||
filter: &RatingFilter,
|
||
) -> Result<usize, dr_catalog::CatalogError> {
|
||
let Some(scope) = scope else {
|
||
return total_images_filtered(catalog, filter);
|
||
};
|
||
|
||
let ids = dr_catalog::collections::descendants(catalog.connection(), scope)?;
|
||
let placeholders = std::iter::repeat_n("?", ids.len())
|
||
.collect::<Vec<_>>()
|
||
.join(",");
|
||
let rated = filter.sql();
|
||
// DISTINCT: an image in both a parent and a child is one photograph, and a
|
||
// count that disagrees with the number of cells drawn is worse than either
|
||
// number alone.
|
||
//
|
||
// Counted through `images` rather than over `collection_members` alone, so
|
||
// `VISIBLE` applies — a trashed photograph is still a member row, and
|
||
// counting it made the header claim images the grid would not draw.
|
||
let sql = format!(
|
||
"SELECT count(DISTINCT i.id) FROM images i
|
||
WHERE {VISIBLE}{rated}
|
||
AND i.id IN (SELECT image_id FROM collection_members
|
||
WHERE collection_id IN ({placeholders}))"
|
||
);
|
||
let params: Vec<rusqlite::types::Value> = ids
|
||
.iter()
|
||
.map(|c| rusqlite::types::Value::Integer(c.0 as i64))
|
||
.collect();
|
||
let n: i64 =
|
||
catalog
|
||
.connection()
|
||
.query_row(&sql, rusqlite::params_from_iter(params.iter()), |r| {
|
||
r.get(0)
|
||
})?;
|
||
Ok(n as usize)
|
||
}
|
||
|
||
/// Total images in the catalog, honouring the rating filter.
|
||
fn total_images_filtered(
|
||
catalog: &Catalog,
|
||
filter: &RatingFilter,
|
||
) -> Result<usize, dr_catalog::CatalogError> {
|
||
let rated = filter.sql();
|
||
let n: i64 = catalog.connection().query_row(
|
||
&format!("SELECT count(*) FROM images i WHERE {VISIBLE}{rated}"),
|
||
[],
|
||
|r| r.get(0),
|
||
)?;
|
||
Ok(n as usize)
|
||
}
|
||
|
||
/// TRACES: FR-CAT-9
|
||
/// How many visible images have their original stored on this device.
|
||
///
|
||
/// Whole-library, like the star counts beside it: the chip says what narrowing
|
||
/// to it would show, so counting only the current window would make it
|
||
/// describe the view it exists to change.
|
||
pub fn local_original_count(catalog: &Catalog) -> Result<usize, dr_catalog::CatalogError> {
|
||
let n: i64 = catalog.connection().query_row(
|
||
&format!(
|
||
"SELECT count(*) FROM images i
|
||
WHERE {VISIBLE}
|
||
AND EXISTS (SELECT 1 FROM image_cache ic
|
||
WHERE ic.image_id = i.id AND ic.tier_actual >= {})",
|
||
dr_types::Tier::Original.stored()
|
||
),
|
||
[],
|
||
|r| r.get(0),
|
||
)?;
|
||
Ok(n as usize)
|
||
}
|
||
|
||
/// Total images in the catalog, unfiltered.
|
||
///
|
||
/// What the scan reports and what the sidebar's "all images" row shows — the
|
||
/// size of the library itself, not of the current view.
|
||
pub fn total_images(catalog: &Catalog) -> Result<usize, dr_catalog::CatalogError> {
|
||
let n: i64 = catalog.connection().query_row(
|
||
&format!("SELECT count(*) FROM images WHERE {VISIBLE_UNALIASED}"),
|
||
[],
|
||
|r| r.get(0),
|
||
)?;
|
||
Ok(n as usize)
|
||
}
|
||
|
||
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 dr_sync::RemoteEntry;
|
||
|
||
#[test]
|
||
fn join_all_preserves_order_regardless_of_completion() {
|
||
// The ordering guarantee is what lets a caller pair results back to
|
||
// their inputs; without it a lane's dates could be attributed to the
|
||
// wrong images.
|
||
let rt = crate::net_runtime::build().unwrap();
|
||
|
||
let out = rt.block_on(async {
|
||
futures_join_all(vec![
|
||
Box::pin(async { 1 }) as std::pin::Pin<Box<dyn std::future::Future<Output = i32>>>,
|
||
Box::pin(async {
|
||
tokio::task::yield_now().await;
|
||
tokio::task::yield_now().await;
|
||
2
|
||
}),
|
||
Box::pin(async {
|
||
tokio::task::yield_now().await;
|
||
3
|
||
}),
|
||
])
|
||
.await
|
||
});
|
||
|
||
assert_eq!(out, vec![1, 2, 3]);
|
||
}
|
||
|
||
#[test]
|
||
fn join_all_of_nothing_completes() {
|
||
let rt = crate::net_runtime::build().unwrap();
|
||
let out: Vec<i32> =
|
||
rt.block_on(async { futures_join_all(Vec::<std::future::Ready<i32>>::new()).await });
|
||
assert!(out.is_empty());
|
||
}
|
||
|
||
#[test]
|
||
fn sweep_lanes_divide_a_chunk_without_loss() {
|
||
// Every image in a chunk must land in exactly one lane: a striding
|
||
// split that dropped or duplicated one would silently under- or
|
||
// double-index the library.
|
||
let chunk: Vec<usize> = (0..SWEEP_CHUNK).collect();
|
||
let lanes: Vec<Vec<usize>> = (0..SWEEP_LANES)
|
||
.map(|l| chunk.iter().skip(l).step_by(SWEEP_LANES).copied().collect())
|
||
.collect();
|
||
|
||
let mut seen: Vec<usize> = lanes.iter().flatten().copied().collect();
|
||
seen.sort_unstable();
|
||
assert_eq!(seen, chunk);
|
||
// Evenly divided, so no lane sits idle while another finishes.
|
||
assert!(lanes.iter().all(|l| l.len() == SWEEP_CHUNK / SWEEP_LANES));
|
||
}
|
||
|
||
#[test]
|
||
fn a_short_chunk_still_covers_every_image() {
|
||
// The last chunk of a library is rarely a full multiple of the lanes.
|
||
let chunk: Vec<usize> = (0..5).collect();
|
||
let lanes: Vec<Vec<usize>> = (0..SWEEP_LANES)
|
||
.map(|l| chunk.iter().skip(l).step_by(SWEEP_LANES).copied().collect())
|
||
.collect();
|
||
let mut seen: Vec<usize> = lanes.iter().flatten().copied().collect();
|
||
seen.sort_unstable();
|
||
assert_eq!(seen, chunk);
|
||
}
|
||
|
||
#[test]
|
||
fn a_scrub_ordinal_matches_the_grid_position() {
|
||
// The scrub's count and the grid's window must use *identical*
|
||
// predicates and ordering, or the view lands somewhere else. Counting
|
||
// only dated images against a grid that also shows undated ones put a
|
||
// click near the end of the axis near the top of the library.
|
||
let catalog = Catalog::in_memory().unwrap();
|
||
let c = catalog.connection();
|
||
c.execute(
|
||
"INSERT INTO roots(id, kind, label) VALUES (1, 'remote', 'lib')",
|
||
[],
|
||
)
|
||
.unwrap();
|
||
|
||
// A mix: dated, undated, and one shadowed by a RAW sibling.
|
||
for (id, name, captured, shadow) in [
|
||
(1i64, "a.CR2", Some(100i64), None),
|
||
(2, "b.CR2", Some(200), None),
|
||
(3, "b.JPG", Some(200), Some(2i64)),
|
||
(4, "c.CR2", Some(300), None),
|
||
(5, "d.CR2", None, None),
|
||
] {
|
||
c.execute(
|
||
"INSERT INTO images(id, root_id, source_ref, captured_at, shadowed_by, added_at)
|
||
VALUES (?1, 1, ?2, ?3, ?4, 0)",
|
||
rusqlite::params![id, name, captured, shadow],
|
||
)
|
||
.unwrap();
|
||
}
|
||
|
||
// The grid's own window, in its own order.
|
||
let cells = read_cells(&catalog, 0, 100).unwrap();
|
||
let names: Vec<&str> = cells.iter().map(|c| c.name.as_str()).collect();
|
||
assert_eq!(
|
||
names,
|
||
vec!["a.CR2", "b.CR2", "c.CR2", "d.CR2"],
|
||
"shadowed hidden, undated last"
|
||
);
|
||
|
||
// Scrubbing to each image's instant must give its index in that list.
|
||
for (when, expected) in [(100i64, 0usize), (200, 1), (300, 2)] {
|
||
let ordinal: i64 = c
|
||
.query_row(
|
||
"SELECT count(*) FROM images
|
||
WHERE shadowed_by IS NULL
|
||
AND captured_at IS NOT NULL
|
||
AND captured_at < ?1",
|
||
[when],
|
||
|r| r.get(0),
|
||
)
|
||
.unwrap();
|
||
assert_eq!(
|
||
ordinal as usize, expected,
|
||
"scrubbing to {when} must land on grid row {expected}"
|
||
);
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn catalog_paths_separate_accounts() {
|
||
// Two accounts on one machine must not share an index, or one
|
||
// library's images appear in the other.
|
||
let a = catalog_path("https://cloud.example", "duncan");
|
||
let b = catalog_path("https://cloud.example", "someone");
|
||
let c = catalog_path("https://other.example", "duncan");
|
||
assert_ne!(a, b);
|
||
assert_ne!(a, c);
|
||
}
|
||
|
||
#[test]
|
||
fn catalog_path_is_filesystem_safe() {
|
||
let p = catalog_path("https://cloud.example.com:8443/nc", "duncan");
|
||
let s = p.to_string_lossy();
|
||
assert!(!s.contains("://"));
|
||
assert!(!s.contains(':') || cfg!(windows));
|
||
}
|
||
|
||
#[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(),
|
||
};
|
||
|
||
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(),
|
||
};
|
||
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(),
|
||
};
|
||
|
||
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(),
|
||
};
|
||
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(),
|
||
};
|
||
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(),
|
||
};
|
||
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(),
|
||
};
|
||
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")
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn cells_are_windowed_not_wholesale() {
|
||
let catalog = Catalog::in_memory().unwrap();
|
||
let images: Vec<RemoteEntry> = (0..50)
|
||
.map(|i| entry(&format!("PhotosRaw/img{i:03}.CR2"), i as u64, 1000))
|
||
.collect();
|
||
let result = dr_sync::ScanResult {
|
||
images,
|
||
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e"))],
|
||
progress: Default::default(),
|
||
};
|
||
persist(&catalog, "PhotosRaw", &result).unwrap();
|
||
|
||
let page = read_cells(&catalog, 10, 5).unwrap();
|
||
assert_eq!(page.len(), 5);
|
||
assert_eq!(page[0].name, "img010.CR2");
|
||
}
|
||
|
||
/// A catalog with `n` images, ready to file into collections.
|
||
fn with_images(n: usize) -> Catalog {
|
||
let catalog = Catalog::in_memory().unwrap();
|
||
let images: Vec<RemoteEntry> = (0..n)
|
||
.map(|i| entry(&format!("PhotosRaw/img{i:03}.CR2"), i as u64, 1000))
|
||
.collect();
|
||
let result = dr_sync::ScanResult {
|
||
images,
|
||
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e"))],
|
||
progress: Default::default(),
|
||
};
|
||
persist(&catalog, "PhotosRaw", &result).unwrap();
|
||
catalog
|
||
}
|
||
|
||
fn image_ids(catalog: &Catalog) -> Vec<dr_types::ImageId> {
|
||
let mut stmt = catalog
|
||
.connection()
|
||
.prepare("SELECT id FROM images ORDER BY source_ref")
|
||
.unwrap();
|
||
stmt.query_map([], |r| Ok(dr_types::ImageId(r.get::<_, i64>(0)? as u64)))
|
||
.unwrap()
|
||
.map(Result::unwrap)
|
||
.collect()
|
||
}
|
||
|
||
#[test]
|
||
fn a_scoped_grid_shows_only_that_collections_images() {
|
||
use dr_catalog::collections::{self as coll, CollectionKind};
|
||
|
||
let catalog = with_images(10);
|
||
let ids = image_ids(&catalog);
|
||
let c = coll::create(
|
||
catalog.connection(),
|
||
"Selects",
|
||
None,
|
||
CollectionKind::Manual,
|
||
)
|
||
.unwrap();
|
||
coll::add_images(catalog.connection(), c, &ids[2..5]).unwrap();
|
||
|
||
let cells = read_cells_scoped(&catalog, Some(c), &RatingFilter::default(), 0, 120).unwrap();
|
||
assert_eq!(cells.len(), 3);
|
||
assert_eq!(
|
||
total_images_scoped(&catalog, Some(c), &RatingFilter::default()).unwrap(),
|
||
3
|
||
);
|
||
// Unscoped is still the whole library.
|
||
assert_eq!(
|
||
total_images_scoped(&catalog, None, &RatingFilter::default()).unwrap(),
|
||
10
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn a_collection_set_shows_its_childrens_images() {
|
||
// A parent whose children hold everything must not read as empty —
|
||
// that is what makes nesting look broken.
|
||
use dr_catalog::collections::{self as coll, CollectionKind};
|
||
|
||
let catalog = with_images(10);
|
||
let ids = image_ids(&catalog);
|
||
let trips =
|
||
coll::create(catalog.connection(), "Trips", None, CollectionKind::Manual).unwrap();
|
||
let iceland = coll::create(
|
||
catalog.connection(),
|
||
"Iceland",
|
||
Some(trips),
|
||
CollectionKind::Manual,
|
||
)
|
||
.unwrap();
|
||
coll::add_images(catalog.connection(), iceland, &ids[0..4]).unwrap();
|
||
|
||
// The parent itself has no direct members at all.
|
||
let cells =
|
||
read_cells_scoped(&catalog, Some(trips), &RatingFilter::default(), 0, 120).unwrap();
|
||
assert_eq!(cells.len(), 4, "the set shows what its children hold");
|
||
assert_eq!(
|
||
total_images_scoped(&catalog, Some(trips), &RatingFilter::default()).unwrap(),
|
||
4
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn an_image_in_both_a_parent_and_a_child_is_shown_once() {
|
||
// The count and the number of cells drawn must agree, or neither is
|
||
// believable.
|
||
use dr_catalog::collections::{self as coll, CollectionKind};
|
||
|
||
let catalog = with_images(10);
|
||
let ids = image_ids(&catalog);
|
||
let trips =
|
||
coll::create(catalog.connection(), "Trips", None, CollectionKind::Manual).unwrap();
|
||
let iceland = coll::create(
|
||
catalog.connection(),
|
||
"Iceland",
|
||
Some(trips),
|
||
CollectionKind::Manual,
|
||
)
|
||
.unwrap();
|
||
coll::add_images(catalog.connection(), trips, &ids[0..2]).unwrap();
|
||
coll::add_images(catalog.connection(), iceland, &ids[0..3]).unwrap();
|
||
|
||
let cells =
|
||
read_cells_scoped(&catalog, Some(trips), &RatingFilter::default(), 0, 120).unwrap();
|
||
assert_eq!(cells.len(), 3, "images 0..3, each once");
|
||
assert_eq!(
|
||
total_images_scoped(&catalog, Some(trips), &RatingFilter::default()).unwrap(),
|
||
3
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn a_scoped_window_still_pages() {
|
||
// FR-CAT-4 applies inside a collection too: a 5,000-image collection
|
||
// must not become 5,000 rows.
|
||
use dr_catalog::collections::{self as coll, CollectionKind};
|
||
|
||
let catalog = with_images(30);
|
||
let ids = image_ids(&catalog);
|
||
let c = coll::create(catalog.connection(), "Big", None, CollectionKind::Manual).unwrap();
|
||
coll::add_images(catalog.connection(), c, &ids).unwrap();
|
||
|
||
let page = read_cells_scoped(&catalog, Some(c), &RatingFilter::default(), 10, 5).unwrap();
|
||
assert_eq!(page.len(), 5);
|
||
assert_eq!(page[0].name, "img010.CR2");
|
||
}
|
||
|
||
#[test]
|
||
fn an_empty_collection_reads_as_empty_rather_than_as_the_whole_library() {
|
||
// The failure that would make scoping useless: an empty IN-list
|
||
// matching everything.
|
||
use dr_catalog::collections::{self as coll, CollectionKind};
|
||
|
||
let catalog = with_images(10);
|
||
let c = coll::create(catalog.connection(), "Empty", None, CollectionKind::Manual).unwrap();
|
||
|
||
assert!(
|
||
read_cells_scoped(&catalog, Some(c), &RatingFilter::default(), 0, 120)
|
||
.unwrap()
|
||
.is_empty()
|
||
);
|
||
assert_eq!(
|
||
total_images_scoped(&catalog, Some(c), &RatingFilter::default()).unwrap(),
|
||
0
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn cells_carry_the_file_id_the_thumbnail_store_keys_on() {
|
||
// Without this the store can never be hit: every launch would refetch
|
||
// every thumbnail over the network.
|
||
let catalog = Catalog::in_memory().unwrap();
|
||
let result = dr_sync::ScanResult {
|
||
images: vec![entry("PhotosRaw/a.CR2", 7777, 30_000_000)],
|
||
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e"))],
|
||
progress: Default::default(),
|
||
};
|
||
persist(&catalog, "PhotosRaw", &result).unwrap();
|
||
|
||
let cells = read_cells(&catalog, 0, 10).unwrap();
|
||
assert_eq!(cells[0].file_id, Some(7777));
|
||
}
|
||
|
||
#[test]
|
||
fn a_stored_thumbnail_survives_a_restart() {
|
||
// The end-to-end property the store exists for: encode, persist,
|
||
// reopen, decode. A second launch must not re-fetch.
|
||
let dir = std::env::temp_dir().join(format!("dr-ui-thumbs-{}", std::process::id()));
|
||
let _ = std::fs::remove_dir_all(&dir);
|
||
|
||
let rgba: Vec<u8> = std::iter::repeat_n([90u8, 140, 200, 255], 32 * 32)
|
||
.flatten()
|
||
.collect();
|
||
{
|
||
let mut store = ThumbStore::open(&dir).unwrap();
|
||
let bytes = dr_thumbs::encode_rgba(32, 32, &rgba).unwrap();
|
||
store
|
||
.put(
|
||
4242,
|
||
dr_thumbs::ThumbSize::Grid,
|
||
&dr_thumbs::Thumbnail {
|
||
width: 32,
|
||
height: 32,
|
||
bytes,
|
||
},
|
||
)
|
||
.unwrap();
|
||
}
|
||
|
||
let store = ThumbStore::open(&dir).unwrap();
|
||
let stored = store
|
||
.get(4242, dr_thumbs::ThumbSize::Grid)
|
||
.unwrap()
|
||
.expect("persisted");
|
||
let (w, h, out) = dr_thumbs::decode_rgba(&stored.bytes).unwrap();
|
||
assert_eq!((w, h), (32, 32));
|
||
// Lossy, so compare approximately — a blue-ish pixel must stay blue.
|
||
assert!(out[2] > out[0], "channel order survived the round trip");
|
||
}
|
||
|
||
#[test]
|
||
fn a_store_hit_is_not_evidence_the_server_is_reachable() {
|
||
// The regression this guards: store hits were delivered as the same
|
||
// `Ready` the network path sends, and the UI took any `Ready` as proof
|
||
// of connectivity. A mostly-cached window then declared "back online"
|
||
// against a server that was down — clearing the banner and kicking off
|
||
// a sweep that immediately failed, on every scroll.
|
||
let dir = std::env::temp_dir().join(format!("dr-ui-provenance-{}", std::process::id()));
|
||
let _ = std::fs::remove_dir_all(&dir);
|
||
|
||
let rgba: Vec<u8> = std::iter::repeat_n([10u8, 20, 30, 255], 8 * 8)
|
||
.flatten()
|
||
.collect();
|
||
let mut store = ThumbStore::open(&dir).unwrap();
|
||
let bytes = dr_thumbs::encode_rgba(8, 8, &rgba).unwrap();
|
||
store
|
||
.put(
|
||
99,
|
||
dr_thumbs::ThumbSize::Grid,
|
||
&dr_thumbs::Thumbnail {
|
||
width: 8,
|
||
height: 8,
|
||
bytes,
|
||
},
|
||
)
|
||
.unwrap();
|
||
|
||
// The split in `spawn_thumbnails`: a hit decodes off local disk and is
|
||
// reported with `from_cache` set, which is what the reachability gate
|
||
// keys on.
|
||
let stored = store
|
||
.get(99, dr_thumbs::ThumbSize::Grid)
|
||
.unwrap()
|
||
.expect("stored");
|
||
let (width, height, rgba) = dr_thumbs::decode_rgba(&stored.bytes).unwrap();
|
||
let hit = ThumbnailReady {
|
||
row: 0,
|
||
width,
|
||
height,
|
||
rgba,
|
||
from_cache: true,
|
||
};
|
||
|
||
assert!(
|
||
hit.from_cache,
|
||
"a thumbnail read from the store must not be mistaken for a fetch"
|
||
);
|
||
|
||
// And the mechanism it feeds: an offline tracker must survive it.
|
||
let mut reach = dr_sync::Reachability::new();
|
||
let now = std::time::Instant::now();
|
||
reach.mark_unreachable("network error".into(), now);
|
||
|
||
if !hit.from_cache {
|
||
reach.mark_reachable(now);
|
||
}
|
||
assert!(
|
||
reach.is_offline(),
|
||
"replaying cached thumbnails must leave offline mode intact"
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn thumbs_live_beside_the_catalog_not_in_the_cache() {
|
||
// They sync to the server and are shared with other clients, so a
|
||
// cache sweep must not discard them.
|
||
let cat = catalog_path("https://cloud.example", "duncan");
|
||
let thumbs = thumbs_dir("https://cloud.example", "duncan");
|
||
assert_eq!(thumbs.parent(), cat.parent());
|
||
}
|
||
|
||
// --- the trash view (FR-CAT-15) ---------------------------------------
|
||
|
||
/// A catalog with `n` scanned images, none trashed.
|
||
fn scanned(n: u64) -> Catalog {
|
||
let catalog = Catalog::in_memory().unwrap();
|
||
let images = (1..=n)
|
||
.map(|i| entry(&format!("PhotosRaw/IMG_{i:04}.CR2"), 1000 + i, 30_000_000))
|
||
.collect();
|
||
persist(
|
||
&catalog,
|
||
"PhotosRaw",
|
||
&dr_sync::ScanResult {
|
||
images,
|
||
directories: vec![(RemotePath::new("PhotosRaw"), dr_sync::Validator::new("e1"))],
|
||
progress: Default::default(),
|
||
},
|
||
)
|
||
.unwrap();
|
||
catalog
|
||
}
|
||
|
||
/// Mark one image trashed, at a given instant.
|
||
fn trash_at(catalog: &Catalog, path: &str, when: i64) {
|
||
let n = catalog
|
||
.connection()
|
||
.execute(
|
||
"UPDATE images SET trashed_at = ?1, trashed_from = source_ref
|
||
WHERE source_ref = ?2",
|
||
rusqlite::params![when, path],
|
||
)
|
||
.unwrap();
|
||
assert_eq!(n, 1, "fixture should have trashed exactly {path}");
|
||
}
|
||
|
||
#[test]
|
||
fn the_trash_lists_what_the_library_hides() {
|
||
// The whole point of the view: these rows exist and no other query in
|
||
// the application will show them.
|
||
let catalog = scanned(3);
|
||
trash_at(&catalog, "PhotosRaw/IMG_0002.CR2", 100);
|
||
|
||
let trashed = read_trashed_cells(&catalog, 0, 50).unwrap();
|
||
assert_eq!(trashed.len(), 1);
|
||
assert_eq!(trashed[0].remote_path, "PhotosRaw/IMG_0002.CR2");
|
||
|
||
// And it has left the library in the same move.
|
||
let live = read_cells_all(&catalog, &RatingFilter::default(), 0, 50).unwrap();
|
||
assert_eq!(live.len(), 2);
|
||
assert!(!live.iter().any(|c| c.remote_path.contains("IMG_0002")));
|
||
}
|
||
|
||
#[test]
|
||
fn an_untrashed_library_has_an_empty_trash() {
|
||
let catalog = scanned(3);
|
||
assert!(read_trashed_cells(&catalog, 0, 50).unwrap().is_empty());
|
||
assert_eq!(total_trashed(&catalog).unwrap(), 0);
|
||
}
|
||
|
||
#[test]
|
||
fn the_trash_count_agrees_with_the_cells_it_lists() {
|
||
// The header and the scrollbar are sized from the count while the grid
|
||
// draws the cells. Two predicates that drift leave the user scrolling
|
||
// through rows that are not there.
|
||
let catalog = scanned(5);
|
||
for (i, when) in [(1, 100), (3, 200), (5, 300)] {
|
||
trash_at(&catalog, &format!("PhotosRaw/IMG_{i:04}.CR2"), when);
|
||
}
|
||
|
||
assert_eq!(total_trashed(&catalog).unwrap(), 3);
|
||
assert_eq!(read_trashed_cells(&catalog, 0, 50).unwrap().len(), 3);
|
||
}
|
||
|
||
#[test]
|
||
fn the_most_recently_trashed_image_is_listed_first() {
|
||
// A mistaken delete is corrected within seconds, so the row the user
|
||
// wants is the one they just made — not the oldest photograph.
|
||
let catalog = scanned(3);
|
||
trash_at(&catalog, "PhotosRaw/IMG_0001.CR2", 100);
|
||
trash_at(&catalog, "PhotosRaw/IMG_0002.CR2", 300);
|
||
trash_at(&catalog, "PhotosRaw/IMG_0003.CR2", 200);
|
||
|
||
let order: Vec<String> = read_trashed_cells(&catalog, 0, 50)
|
||
.unwrap()
|
||
.into_iter()
|
||
.map(|c| c.remote_path)
|
||
.collect();
|
||
assert_eq!(
|
||
order,
|
||
vec![
|
||
"PhotosRaw/IMG_0002.CR2".to_string(),
|
||
"PhotosRaw/IMG_0003.CR2".to_string(),
|
||
"PhotosRaw/IMG_0001.CR2".to_string(),
|
||
]
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn a_shadowed_jpeg_is_not_listed_beside_the_raw_it_belongs_to() {
|
||
// The inversion applies to `trashed_at` only. Trashing a RAW takes its
|
||
// sibling JPEG with it, and listing both would offer to restore the
|
||
// same frame twice.
|
||
let catalog = scanned(2);
|
||
let c = catalog.connection();
|
||
let raw: i64 = c
|
||
.query_row(
|
||
"SELECT id FROM images WHERE source_ref = 'PhotosRaw/IMG_0001.CR2'",
|
||
[],
|
||
|r| r.get(0),
|
||
)
|
||
.unwrap();
|
||
c.execute(
|
||
"UPDATE images SET shadowed_by = ?1 WHERE source_ref = 'PhotosRaw/IMG_0002.CR2'",
|
||
[raw],
|
||
)
|
||
.unwrap();
|
||
|
||
trash_at(&catalog, "PhotosRaw/IMG_0001.CR2", 100);
|
||
trash_at(&catalog, "PhotosRaw/IMG_0002.CR2", 100);
|
||
|
||
let trashed = read_trashed_cells(&catalog, 0, 50).unwrap();
|
||
assert_eq!(trashed.len(), 1, "one photograph, not two");
|
||
assert_eq!(trashed[0].remote_path, "PhotosRaw/IMG_0001.CR2");
|
||
assert_eq!(total_trashed(&catalog).unwrap(), 1, "and the count agrees");
|
||
}
|
||
|
||
#[test]
|
||
fn the_trash_window_pages_like_the_grid_does() {
|
||
// The trash uses the same windowed read as the library, so a large one
|
||
// must not try to draw itself in a single query.
|
||
let catalog = scanned(6);
|
||
for i in 1..=6 {
|
||
trash_at(
|
||
&catalog,
|
||
&format!("PhotosRaw/IMG_{i:04}.CR2"),
|
||
100 + i as i64,
|
||
);
|
||
}
|
||
|
||
let first = read_trashed_cells(&catalog, 0, 2).unwrap();
|
||
let second = read_trashed_cells(&catalog, 2, 2).unwrap();
|
||
assert_eq!(first.len(), 2);
|
||
assert_eq!(second.len(), 2);
|
||
assert!(
|
||
first
|
||
.iter()
|
||
.all(|a| !second.iter().any(|b| b.image_id == a.image_id)),
|
||
"pages must not overlap"
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn a_rating_filter_does_not_hide_anything_in_the_trash() {
|
||
// The filter narrows a culling pass. A trash that dropped rows because
|
||
// of a filter set elsewhere would look like it had lost them.
|
||
let catalog = scanned(3);
|
||
trash_at(&catalog, "PhotosRaw/IMG_0001.CR2", 100);
|
||
|
||
// Nothing here is rated, so a four-star library filter would empty any
|
||
// view that honoured it.
|
||
let strict = RatingFilter {
|
||
min_rating: 4,
|
||
..Default::default()
|
||
};
|
||
assert!(read_cells_all(&catalog, &strict, 0, 50).unwrap().is_empty());
|
||
assert_eq!(read_trashed_cells(&catalog, 0, 50).unwrap().len(), 1);
|
||
}
|
||
|
||
fn entry(path: &str, file_id: u64, size: u64) -> RemoteEntry {
|
||
RemoteEntry {
|
||
id: RemoteId::Stable(file_id),
|
||
path: RemotePath::new(path),
|
||
kind: dr_sync::EntryKind::File,
|
||
validator: dr_sync::Validator::new("v"),
|
||
size,
|
||
modified: None,
|
||
has_preview: false,
|
||
}
|
||
}
|
||
|
||
// --- the local-only filter (FR-CAT-9) ---------------------------------
|
||
|
||
/// Record that an image's original is held locally at `tier`.
|
||
fn cache_at(catalog: &Catalog, id: dr_types::ImageId, tier: dr_types::Tier) {
|
||
catalog
|
||
.connection()
|
||
.execute(
|
||
"INSERT INTO image_cache (image_id, tier_actual, bytes)
|
||
VALUES (?1, ?2, 0)",
|
||
rusqlite::params![id.0 as i64, tier.stored()],
|
||
)
|
||
.unwrap();
|
||
}
|
||
|
||
#[test]
|
||
fn local_only_shows_just_the_images_held_here() {
|
||
let catalog = with_images(10);
|
||
let ids = image_ids(&catalog);
|
||
for id in &ids[0..3] {
|
||
cache_at(&catalog, *id, dr_types::Tier::Original);
|
||
}
|
||
|
||
let filter = RatingFilter {
|
||
local_only: true,
|
||
..Default::default()
|
||
};
|
||
let cells = read_cells_all(&catalog, &filter, 0, 120).unwrap();
|
||
assert_eq!(cells.len(), 3);
|
||
// The count the header shows must agree with the cells drawn, which is
|
||
// the whole reason the predicate lives in SQL rather than in a
|
||
// post-filter over the rows.
|
||
assert_eq!(total_images_filtered(&catalog, &filter).unwrap(), 3);
|
||
assert_eq!(local_original_count(&catalog).unwrap(), 3);
|
||
}
|
||
|
||
#[test]
|
||
fn a_cached_preview_is_not_a_local_original() {
|
||
// The filter answers "can I open this in develop right now", and a
|
||
// preview cannot. Counting it would put images in the offline set that
|
||
// fail the moment they are clicked.
|
||
let catalog = with_images(5);
|
||
let ids = image_ids(&catalog);
|
||
cache_at(&catalog, ids[0], dr_types::Tier::Preview);
|
||
cache_at(&catalog, ids[1], dr_types::Tier::Original);
|
||
|
||
let filter = RatingFilter {
|
||
local_only: true,
|
||
..Default::default()
|
||
};
|
||
assert_eq!(read_cells_all(&catalog, &filter, 0, 120).unwrap().len(), 1);
|
||
assert_eq!(local_original_count(&catalog).unwrap(), 1);
|
||
}
|
||
|
||
#[test]
|
||
fn local_only_composes_with_the_rating_filter() {
|
||
// "Five-star frames I can actually edit on this train" is one filter,
|
||
// not a mode that replaces the others.
|
||
let catalog = with_images(6);
|
||
let ids = image_ids(&catalog);
|
||
for id in &ids[0..4] {
|
||
cache_at(&catalog, *id, dr_types::Tier::Original);
|
||
}
|
||
// Rate two of the cached ones, and one that is not cached.
|
||
for id in [ids[0], ids[1], ids[5]] {
|
||
dr_catalog::rating::set_rating(catalog.connection(), id, 5).unwrap();
|
||
}
|
||
|
||
let filter = RatingFilter {
|
||
min_rating: 5,
|
||
local_only: true,
|
||
..Default::default()
|
||
};
|
||
let cells = read_cells_all(&catalog, &filter, 0, 120).unwrap();
|
||
assert_eq!(cells.len(), 2, "five-starred AND held locally");
|
||
assert_eq!(total_images_filtered(&catalog, &filter).unwrap(), 2);
|
||
}
|
||
|
||
#[test]
|
||
fn an_empty_cache_is_not_an_empty_library() {
|
||
// The unfiltered grid must not depend on the cache table having rows —
|
||
// a library nothing has been downloaded from is still a full library.
|
||
let catalog = with_images(4);
|
||
assert_eq!(local_original_count(&catalog).unwrap(), 0);
|
||
assert_eq!(
|
||
read_cells_all(&catalog, &RatingFilter::default(), 0, 120)
|
||
.unwrap()
|
||
.len(),
|
||
4
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn local_only_counts_as_a_narrowing_filter() {
|
||
// `is_unfiltered` gates the "filtered" indicator. Reporting this one as
|
||
// unfiltered would leave a narrowed grid looking like the whole
|
||
// library, which is the state the indicator exists to prevent.
|
||
assert!(RatingFilter::default().is_unfiltered());
|
||
assert!(!RatingFilter {
|
||
local_only: true,
|
||
..Default::default()
|
||
}
|
||
.is_unfiltered());
|
||
}
|
||
}
|