Files
DarkRoom/ui/dr-ui/src/derived_sync.rs
T
dtourolleandClaude Opus 5 9e47133304
Build and test / android-image (push) Canceled after 1m38s
Build and test / Android (aarch64) (push) Canceled after 0s
Build and test / Desktop (Linux) (push) Canceled after 1m34s
Build and test / Layer separation (push) Canceled after 0s
🐳 Android image / Build and push (push) Canceled after 1m38s
Traceability / Requirement traces (push) Successful in 1m2s
Run the formatter over the shard-naming change
`cargo fmt --all -- --check` is a CI gate and 9a51cc8 landed two files past
it, so the build has been red on master since regardless of what came after.
Whitespace only — a wrapped `Ok(...)` and two `assert_eq!`s split across
lines. No logic is touched and the tests are unchanged.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-16 22:29:28 +02:00

468 lines
18 KiB
Rust

//! TRACES: FR-CAT-3 | FR-CAT-7 | FR-NC-7
//! Pushing derived state to Nextcloud: thumbnail shards and the catalog.
//!
//! # What travels, and why only this
//!
//! Sidecars are handled elsewhere ([`crate::library::spawn_sidecar_writes`])
//! and are the *authoritative* store — they are the reason a catalog can be
//! deleted and rebuilt (ARCH §6.12). What moves here is derived state that is
//! merely expensive:
//!
//! - **Thumbnail shards.** A thumbnail costs a range fetch plus a decode, and
//! is byte-identical for every client looking at the same file. A second
//! device that downloads the shards gets a full grid without touching a
//! single RAW — hours of indexing against a few hundred MB of transfer.
//! - **The catalog**, for its collections. Every other thing the catalog holds
//! has authoritative backing in a sidecar; a manually assembled collection
//! does not, so without this it exists on one machine only.
//!
//! # Why sealed shards make this cheap
//!
//! A shard stops being written once it reaches its cap, and is never rewritten
//! after — deleting a thumbnail tombstones it in the index rather than editing
//! the sealed blob. So a client that has downloaded a sealed shard never needs
//! to ask about it again, and an up-to-date client transfers only the index and
//! whichever shard is currently open. That is the whole reason for sharding at
//! 25 MB rather than keeping one growing file.
//!
//! # Where it lives
//!
//! Under the library root, in a dotted folder beside the trash. The root is the
//! only place the user granted access to, and writing outside it may cross a
//! share boundary the account cannot write to. The scanner excludes it by the
//! same mechanism that excludes the trash.
use std::path::{Path, PathBuf};
use dr_sync::{RemoteBackend, RemoteId, RemotePath};
use dr_sync_nextcloud::{AppCredentials, NextcloudBackend};
use dr_thumbs::ThumbStore;
/// Folder under the library root holding derived state.
///
/// Defined by the scanner, which must exclude it: a walk that indexed this
/// folder would pay a listing for it on every sync of every device.
pub use dr_sync::scan::DERIVED_DIR;
/// What a sync pass did, for logging and for telling the user.
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct SyncReport {
pub shards_uploaded: usize,
pub shards_downloaded: usize,
pub thumbnails_adopted: usize,
pub catalog_uploaded: bool,
pub catalog_merged: bool,
pub collections_gained: usize,
}
impl SyncReport {
pub fn did_anything(&self) -> bool {
self.shards_uploaded > 0
|| self.shards_downloaded > 0
|| self.catalog_uploaded
|| self.catalog_merged
}
}
/// Progress from the sync worker.
#[derive(Debug)]
pub enum SyncMessage {
Status(String),
Finished(Box<SyncReport>),
Failed(String),
}
/// Push shards and the catalog, and take anything newer from the server.
///
/// Runs on its own thread with its own runtime, like every other network path
/// here — the Slint loop must never block (NFR-P9).
pub fn spawn_sync(
creds: AppCredentials,
user_id: String,
root: String,
thumbs_dir: PathBuf,
catalog_path: PathBuf,
scratch: PathBuf,
) -> std::sync::mpsc::Receiver<SyncMessage> {
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
// A current-thread runtime here is what made the library scan itself
// rather than adopt the shards the server already had; see net_runtime.
let rt = match crate::net_runtime::build() {
Ok(rt) => rt,
Err(e) => {
let _ = tx.send(SyncMessage::Failed(e.to_string()));
return;
}
};
rt.block_on(async {
let backend = match NextcloudBackend::new(&creds, &user_id) {
Ok(b) => b,
Err(e) => {
let _ = tx.send(SyncMessage::Failed(e.to_string()));
return;
}
};
match run(&backend, &root, &thumbs_dir, &catalog_path, &scratch, &tx).await {
Ok(report) => {
let _ = tx.send(SyncMessage::Finished(Box::new(report)));
}
Err(e) => {
let _ = tx.send(SyncMessage::Failed(e));
}
}
});
});
rx
}
async fn run(
backend: &NextcloudBackend,
root: &str,
thumbs_dir: &Path,
catalog_path: &Path,
scratch: &Path,
tx: &std::sync::mpsc::Sender<SyncMessage>,
) -> Result<SyncReport, String> {
let mut report = SyncReport::default();
let base = derived_path(root);
// The folder may not exist on a first sync. Creating it unconditionally is
// cheaper than probing, and an existing folder is not an error.
let _ = backend.create_dir(&base).await;
let _ = tx.send(SyncMessage::Status("checking thumbnails…".into()));
sync_shards(backend, &base, thumbs_dir, scratch, &mut report).await?;
let _ = tx.send(SyncMessage::Status("checking collections…".into()));
sync_catalog(backend, &base, catalog_path, scratch, &mut report).await?;
Ok(report)
}
/// The derived folder for a library root.
fn derived_path(root: &str) -> RemotePath {
if root.is_empty() {
RemotePath::new(DERIVED_DIR)
} else {
RemotePath::new(format!("{root}/{DERIVED_DIR}"))
}
}
/// Exchange thumbnail shards with the server.
///
/// Upload what the server lacks, download what we lack. Sealed shards are
/// immutable, so a name match is a content match and nothing needs comparing
/// beyond existence — which is what keeps a steady-state sync to one listing.
///
/// # Why the name carries a client
///
/// Shard ids are per-store: every client fills its own numbering from 0, so
/// "shard 3" names different thumbnails on each device. A flat `shard-0003`
/// remote namespace therefore has two clients writing one name — the second
/// upload overwrites content the first still believes is published — and
/// leaves a client no way to tell a peer's shard 3 from its own, so the only
/// safe reading of "I already have 3" is to skip it and never adopt anything.
/// Qualifying the name with [`ThumbStore::client_id`] gives each store its own
/// namespace, and [`ThumbStore::adopted`] then tracks what has been merged by
/// remote name instead of by our own ids.
async fn sync_shards(
backend: &NextcloudBackend,
base: &RemotePath,
thumbs_dir: &Path,
scratch: &Path,
report: &mut SyncReport,
) -> Result<(), String> {
let store = match ThumbStore::open(thumbs_dir) {
Ok(s) => s,
Err(e) => {
// A store that will not open can be neither read nor merged into,
// so there is no half of this worth attempting. It is not a sync
// failure: the next pass retries once the store is openable.
log::debug!("thumbnail store unavailable: {e}");
return Ok(());
}
};
let client = store.client_id().to_string();
let remote: std::collections::HashMap<String, u64> = backend
.list(base, None)
.await
.map(|entries| {
entries
.into_iter()
.filter(|e| e.kind == dr_sync::EntryKind::File)
.map(|e| (e.path.name().to_string(), e.size))
.collect()
})
// A missing folder lists as an error on some servers; treat it as empty
// rather than aborting a first sync.
.unwrap_or_default();
let local = store.shards().map_err(|e| e.to_string())?;
// ---- upload ----------------------------------------------------------
for shard in &local {
let path = store.shard_path(shard.id);
let Ok(bytes) = std::fs::read(&path) else {
continue;
};
let name = shard_name(&client, shard.id);
// A sealed shard the server already has is byte-identical by
// construction, so its presence is proof enough — and with the client
// in the name, no one else can have written it. The open shard is
// re-uploaded whenever its size differs, which is the only way it
// changes.
let skip = match remote.get(&name) {
Some(_) if shard.sealed => true,
Some(size) => *size == bytes.len() as u64,
None => false,
};
if skip {
continue;
}
let target = RemotePath::new(format!("{}/{name}", base.as_str()));
match backend.put(&target, bytes, None).await {
Ok(_) => report.shards_uploaded += 1,
// One shard failing must not abort the rest: they are independent
// and the next pass retries.
Err(e) => log::warn!("uploading {name}: {e}"),
}
}
// ---- download --------------------------------------------------------
let mut store = store;
for (name, size) in &remote {
let Some((owner, id)) = parse_shard(name) else {
continue;
};
if owner == client {
continue;
}
// Merged already, at the size it still has. A sealed shard never
// reaches here twice; a peer's open one does each time it grows, which
// is what carries its later thumbnails across.
if store.adopted(name) == Some(*size) {
continue;
}
if owner.is_empty() && legacy_upload_of_ours(&store, id, *size) {
let _ = store.record_adopted(name, *size);
continue;
}
let source = RemotePath::new(format!("{}/{name}", base.as_str()));
let bytes = match backend.get(&RemoteId::Path(source), None).await {
Ok(b) => b,
Err(e) => {
log::warn!("downloading {name}: {e}");
continue;
}
};
// Written to scratch and merged, rather than dropped into the store
// directory: a downloaded shard's *id* is the other device's numbering,
// and two devices independently fill shard 0.
let tmp = scratch.join(name);
if std::fs::write(&tmp, &bytes).is_err() {
continue;
}
match store.merge_shard(&tmp) {
Ok(n) => {
report.shards_downloaded += 1;
report.thumbnails_adopted += n;
// Recorded only on success, so a failed merge is retried next
// pass rather than written off.
let _ = store.record_adopted(name, *size);
}
Err(e) => log::warn!("merging {name}: {e}"),
}
let _ = std::fs::remove_file(&tmp);
}
Ok(())
}
/// Whether a flat-named remote shard is this client's own earlier upload.
///
/// Before the name carried a client every client wrote `shard-NNNN.sqlite`, so
/// the folder still holds files with nothing in the name to say whose they
/// are. A local shard of the same id and the same size is ours by
/// construction — the same identity argument the upload path makes for
/// skipping a sealed shard the server already has — and skipping those is what
/// keeps the rename from costing every client a re-download of its whole
/// store. Being wrong costs a peer's shard going unmerged and its thumbnails
/// being derived locally instead; it loses nothing, and two independently
/// filled 25 MB databases landing on the same byte count is not a real case.
fn legacy_upload_of_ours(store: &ThumbStore, id: u32, remote_size: u64) -> bool {
std::fs::metadata(store.shard_path(id))
.map(|m| m.len() == remote_size)
.unwrap_or(false)
}
/// Exchange the catalog, for its collections.
///
/// Only collections merge — see [`dr_catalog::sync`]. The rest of a catalog
/// describes local state (folder ETags, cache paths, job rows) and importing
/// another device's version would be actively wrong.
async fn sync_catalog(
backend: &NextcloudBackend,
base: &RemotePath,
catalog_path: &Path,
scratch: &Path,
report: &mut SyncReport,
) -> Result<(), String> {
let remote_name = "catalog.sqlite";
let target = RemotePath::new(format!("{}/{remote_name}", base.as_str()));
// ---- take theirs first -----------------------------------------------
//
// Merging before uploading means our upload carries the union rather than
// only our own half, so a third device syncing next gets everything in one
// fetch.
if let Ok(bytes) = backend.get(&RemoteId::Path(target.clone()), None).await {
let downloaded = scratch.join("catalog-remote.sqlite");
if std::fs::write(&downloaded, &bytes).is_ok() {
match dr_catalog::Catalog::open(catalog_path) {
Ok(catalog) => match catalog.merge_remote_catalog(&downloaded) {
Ok(merge) => {
report.catalog_merged = true;
report.collections_gained = merge.inserted + merge.updated;
}
Err(e) => log::warn!("merging remote catalog: {e}"),
},
Err(e) => log::warn!("opening catalog to merge: {e}"),
}
let _ = std::fs::remove_file(&downloaded);
}
}
// ---- then push ours --------------------------------------------------
//
// Never the live file: committed transactions can sit in the `-wal` with
// the main file lagging, so copying it uploads a torn snapshot. The backup
// API serialises against writers instead of racing them.
let snapshot = scratch.join("catalog-upload.sqlite");
let catalog = dr_catalog::Catalog::open(catalog_path).map_err(|e| e.to_string())?;
catalog
.snapshot_for_upload(&snapshot)
.map_err(|e| e.to_string())?;
let bytes = std::fs::read(&snapshot).map_err(|e| e.to_string())?;
match backend.put(&target, bytes, None).await {
Ok(_) => report.catalog_uploaded = true,
Err(e) => log::warn!("uploading catalog: {e}"),
}
let _ = std::fs::remove_file(&snapshot);
Ok(())
}
fn shard_name(client: &str, id: u32) -> String {
format!("shard-{client}-{id:04}.sqlite")
}
/// The client that wrote a remote shard and its id in that client's numbering,
/// or `None` if the name is not a shard.
///
/// Also accepts the flat `shard-NNNN.sqlite` written before names carried a
/// client, reporting an empty owner: those belong to nobody identifiable, so
/// they read as foreign and are adopted once like any peer's. Nothing is ever
/// uploaded under that form again.
///
/// Guards the download loop against adopting the catalog, a stray file, or
/// anything else the folder happens to contain.
fn parse_shard(name: &str) -> Option<(&str, u32)> {
let stem = name.strip_prefix("shard-")?.strip_suffix(".sqlite")?;
match stem.rsplit_once('-') {
Some((client, id)) => Some((client, id.parse().ok()?)),
None => Some(("", stem.parse().ok()?)),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn derived_folder_sits_under_the_library_root() {
// Outside the root the account may not have write access — the root is
// the only thing the user granted.
assert_eq!(
derived_path("PhotosRaw").as_str(),
"PhotosRaw/.darkroom-derived"
);
// A library at the account root still gets a relative path.
assert_eq!(derived_path("").as_str(), ".darkroom-derived");
}
#[test]
fn shard_names_round_trip() {
assert_eq!(
shard_name("a1b2c3d4e5f6", 0),
"shard-a1b2c3d4e5f6-0000.sqlite"
);
assert_eq!(
shard_name("a1b2c3d4e5f6", 42),
"shard-a1b2c3d4e5f6-0042.sqlite"
);
assert_eq!(
parse_shard(&shard_name("a1b2c3d4e5f6", 7)),
Some(("a1b2c3d4e5f6", 7))
);
}
#[test]
fn two_clients_shard_three_are_different_files() {
// The whole point: one client's numbering must not name another's
// shard, or the second upload overwrites the first's content and
// neither can tell the other's shards from its own.
assert_ne!(shard_name("aaaa", 3), shard_name("bbbb", 3));
assert_eq!(parse_shard(&shard_name("aaaa", 3)).unwrap().0, "aaaa");
assert_eq!(parse_shard(&shard_name("bbbb", 3)).unwrap().0, "bbbb");
}
#[test]
fn flat_names_read_as_belonging_to_nobody() {
// Written before the name carried a client. They must still parse, so
// a library synced by an older build is not stranded, and they must
// not match any live client id, so they are never mistaken for ours.
assert_eq!(parse_shard("shard-0042.sqlite"), Some(("", 42)));
assert_ne!(parse_shard("shard-0042.sqlite").unwrap().0, "a1b2c3d4e5f6");
}
#[test]
fn non_shard_files_are_not_adopted() {
// The folder also holds the catalog; downloading it as a shard would
// hand a catalog to the thumbnail merger.
assert_eq!(parse_shard("catalog.sqlite"), None);
assert_eq!(parse_shard("shard-0000.sqlite-wal"), None);
assert_eq!(parse_shard("notes.txt"), None);
assert_eq!(parse_shard("shard-abc.sqlite"), None);
assert_eq!(parse_shard("shard-a1b2c3-notanid.sqlite"), None);
}
#[test]
fn a_report_that_did_nothing_says_so() {
assert!(!SyncReport::default().did_anything());
assert!(SyncReport {
shards_uploaded: 1,
..Default::default()
}
.did_anything());
// Adopting thumbnails without moving a shard cannot happen, but the
// report must not claim work on collections alone either.
assert!(SyncReport {
catalog_merged: true,
..Default::default()
}
.did_anything());
}
}