//! 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; 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, /// Membership rows the merge brought in (FR-CAT-7). /// /// Counted apart from `collections_gained`, which counts the collections /// themselves. A sync that files 127 photographs into a collection both /// devices already had gains *no* collection, and reporting only the /// former left the sidebar showing no count beside a full collection. pub members_gained: usize, // Face data is counted apart from thumbnails for the same reason keywords // are counted apart from collections: "adopted 4,812 faces" is a sentence // the user can act on, and folding it into the thumbnail count would hide // the one number that says whether this device still has hours of indexing // ahead of it. pub face_shards_uploaded: usize, pub face_shards_downloaded: usize, /// Images whose faces this device took from a peer instead of detecting. pub faces_adopted: usize, } impl SyncReport { pub fn did_anything(&self) -> bool { self.shards_uploaded > 0 || self.shards_downloaded > 0 || self.catalog_uploaded || self.catalog_merged || self.face_shards_uploaded > 0 || self.face_shards_downloaded > 0 } } /// Progress from the sync worker. #[derive(Debug)] pub enum SyncMessage { Status(String), Finished(Box), 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 { 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 crate::remote::connect(&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: &dyn RemoteBackend, root: &str, thumbs_dir: &Path, catalog_path: &Path, scratch: &Path, tx: &std::sync::mpsc::Sender, ) -> Result { 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 faces…".into())); sync_face_shards(backend, &base, catalog_path, scratch, &mut report, tx).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: &dyn RemoteBackend, 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 = 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(()) } /// Push and pull face shards, so a second device does not re-index the library. /// /// Mirrors [`sync_shards`] deliberately, down to the sealed-shard skip and the /// adopted ledger: face data has exactly the properties that made that design /// right for thumbnails. It is bulk, it is immutable once written, and it is /// byte-identical on every device, because the same model over the same proxy /// is deterministic. /// /// The catalog is the source and the destination; the shards are only the /// carrier. So this exports the catalog's new faces into the local shard store /// first, syncs the shards, and imports whatever arrived back into the catalog. async fn sync_face_shards( backend: &dyn RemoteBackend, base: &RemotePath, catalog_path: &Path, scratch: &Path, report: &mut SyncReport, tx: &std::sync::mpsc::Sender, ) -> Result<(), String> { use dr_catalog::face_shard::{self, FaceShardStore}; // Beside the catalog, next to the thumbnails, and under the same folder the // scanner already excludes. let dir = match catalog_path.parent() { Some(p) => p.join("faces"), None => return Ok(()), }; let mut store = match FaceShardStore::open(&dir) { Ok(s) => s, Err(e) => { // Not a sync failure: the next pass retries once it opens. log::debug!("face shard store unavailable: {e}"); return Ok(()); } }; let client = store.client_id().to_string(); let model = crate::identity_ui::MODEL_ID; // ---- everything this device has detected, into the shards ------------ if let Ok(catalog) = dr_catalog::Catalog::open(catalog_path) { // The first export after a re-index walks the whole library and writes // thousands of images. Saying so is the difference between a sync that // looks slow and one that looks broken — this pass sat on "checking // faces…" for eight minutes and was reported as a hang. let mut say = |done: usize, total: usize| { let _ = tx.send(SyncMessage::Status(format!( "preparing faces to send: {done}/{total}" ))); }; match face_shard::export_to_shards_reporting( catalog.connection(), &mut store, model, &mut say, ) { Ok(0) => {} Ok(n) => log::info!("face sync: {n} newly indexed image(s) ready to upload"), Err(e) => log::warn!("face sync: exporting to shards: {e}"), } } // Its own folder under the derived directory, so a client that does not // care about faces lists thumbnails without paging past them. let face_base = RemotePath::new(format!("{}/faces", base.as_str())); let _ = backend.create_dir(&face_base).await; let remote: std::collections::HashMap = backend .list(&face_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() }) .unwrap_or_default(); // ---- upload ---------------------------------------------------------- // // Every commit since the last pass is still in a write-ahead log, and a // shard is uploaded by reading its file — so without this the upload would // ship a database missing precisely the faces just exported. if let Err(e) = store.checkpoint() { log::warn!("face sync: checkpointing the shard store: {e}"); } let local = store.shards().map_err(|e| e.to_string())?; for (n, shard) in local.iter().enumerate() { let name = shard_name(&client, shard.id); let path = store.shard_path(shard.id); // **Decided before the file is read.** A sealed shard the server // already has is byte-identical by construction, and the client is in // the name so nobody else could have written it — the name alone // settles it. Reading first meant every idle sync pulled ninety-four // megabytes off disk to conclude it had nothing to send. let on_disk = std::fs::metadata(&path).map(|m| m.len()).unwrap_or(0); let skip = match remote.get(&name) { Some(_) if shard.sealed => true, Some(size) => *size == on_disk, None => false, }; if skip { continue; } let Ok(bytes) = std::fs::read(&path) else { continue; }; // Face shards carry crops and run to tens of megabytes each, so a // single one is a visible wait on any connection. Announced after the // skip, or an idle pass claims to be sending five shards and sends // none; and before the put, because the wait is the upload. let _ = tx.send(SyncMessage::Status(format!( "sending faces: shard {}/{} ({} MB)", n + 1, local.len(), bytes.len() / 1_048_576 ))); let target = RemotePath::new(format!("{}/{name}", face_base.as_str())); match backend.put(&target, bytes, None).await { Ok(_) => report.face_shards_uploaded += 1, Err(e) => log::warn!("uploading face shard {name}: {e}"), } } // ---- download -------------------------------------------------------- for (name, size) in &remote { let Some((owner, _)) = parse_shard(name) else { continue; }; if owner == client || store.has_adopted(name, *size) { continue; } // Taking in a peer's shard means a download and then a row-by-row // merge, both of which take a while on a shard carrying crops. let _ = tx.send(SyncMessage::Status(format!( "taking in faces from another device ({} MB)", size / 1_048_576 ))); let source = RemotePath::new(format!("{}/{name}", face_base.as_str())); let bytes = match backend.get(&RemoteId::Path(source), None).await { Ok(b) => b, Err(e) => { log::warn!("downloading face shard {name}: {e}"); continue; } }; // Into scratch and merged, never 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(_) => { report.face_shards_downloaded += 1; // Recorded only on success, so a failed merge is retried next // pass rather than written off. let _ = store.mark_adopted(name, *size); } Err(e) => log::warn!("merging face shard {name}: {e}"), } let _ = std::fs::remove_file(&tmp); } // ---- and back into the catalog --------------------------------------- // // Last, and unconditionally rather than only when something downloaded: a // previous pass may have merged shards into the store and then failed // before importing, and this is what recovers from that. if let Ok(catalog) = dr_catalog::Catalog::open(catalog_path) { match face_shard::import_from_shards(catalog.connection(), &store, model) { Ok(0) => {} Ok(n) => { report.faces_adopted = n; log::info!("face sync: adopted {n} image(s) already indexed elsewhere"); } Err(e) => log::warn!("face sync: importing from shards: {e}"), } } 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: &dyn RemoteBackend, 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; report.members_gained = merge.members_added; } 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()); } }