//! TRACES: FR-CAT-3 | FR-CAT-7 | FR-NC-7 //! Pushing derived state to the library: 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::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath}; 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( conn: Connection, 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(&conn) { 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); // **Size, sealed or not.** This used to take the server merely *having* // a sealed shard as proof it had the whole thing — "byte-identical by // construction". It is not: an open shard is uploaded on every pass as // it grows, so the server routinely holds a partial copy of a shard // that is later filled and sealed. From the moment of sealing, the name // matched and the size was never looked at again, and the completed // shard could never be sent. The client id in the name still means // nobody else can have written it, so size is a sound comparison. let skip = match remote.get(&name) { 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())); // Fetched where it is only a placeholder: a shard that will not open // is a peer's thumbnails never merging, and on a library the client // keeps dehydrated that would be every shard, every pass, silently. let bytes = match read_derived(backend, &source).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**, and on size rather than mere // presence. Reading first meant every idle sync pulled ninety-four // megabytes off disk to conclude it had nothing to send. // // The presence test was worse than wasteful. A sealed shard was taken // to be byte-identical to whatever the server already had under that // name — but an *open* shard is uploaded on every pass as it fills, so // the server ordinarily holds a partial copy of a shard that is later // completed and sealed. Sealing then froze that partial copy in place: // the name matched, the size was never consulted, and the finished // shard was skipped for ever. This library's shard 0 sat on the server // at 2.7 MB against 20 MB on disk, and the tablet adopted the 2.7 MB — // which is why it showed a fraction of the faces and never caught up. let on_disk = std::fs::metadata(&path).map(|m| m.len()).unwrap_or(0); let skip = match remote.get(&name) { 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())); // Fetched where it is only a placeholder, for the reason the thumbnail // shards are: otherwise a peer's faces never arrive and nothing says so. let bytes = match read_derived(backend, &source).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. // TRACES: FR-NC-9 | FR-NC-6c // A read that fails for any reason other than "there is not one yet" must // stop the upload below. This is a read-modify-write over a file another // device also writes, so skipping the read does not merely lose an // optimisation — it turns the write into a clobber, and the other device's // collections and their members go with it. // // The shape was previously `if let Ok(bytes) = ...`, which swallowed every // failure into "no remote catalog" and carried straight on to the upload. let theirs = match read_derived(backend, &target).await { Ok(bytes) => Some(bytes), // Genuinely the first sync of this library. Nothing to merge, and // ours is the whole truth. Err(RemoteError::NotFound(_)) => None, Err(e) => { log::warn!( "not pushing the catalog: the copy on the server could not be read ({e}); \ uploading over it would discard whatever another device put there" ); return Ok(()); } }; if let Some(bytes) = theirs { 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; } // Unreadable is not the same as absent: it may be a newer // format, or a torn upload. Ours must not go over it. Err(e) => { log::warn!("not pushing the catalog: merging the server's copy: {e}"); let _ = std::fs::remove_file(&downloaded); return Ok(()); } }, Err(e) => { log::warn!("not pushing the catalog: opening ours to merge: {e}"); let _ = std::fs::remove_file(&downloaded); return Ok(()); } } 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(()) } /// TRACES: FR-NC-6c /// Read a derived file, fetching its content first if only a placeholder is /// here. /// /// Derived state lives *inside the library folder*, so on a placeholder /// library a sync client dehydrates a shard or a catalog snapshot exactly as /// it dehydrates a photograph. Unlike a photograph, these are ours, and none of /// them can be skipped: a shard that will not open is face data that never /// merges, and a catalog snapshot that will not open is the other device's /// collections. /// /// So this fetches rather than giving up — and where it cannot, it says so /// with the error rather than an empty result, because the callers below treat /// "nothing there" as licence to write their own copy (ARCH §9.0a). async fn read_derived( backend: &dyn RemoteBackend, path: &RemotePath, ) -> Result, RemoteError> { let id = RemoteId::Path(path.clone()); match backend.get(&id, None).await { Err(RemoteError::NotMaterialised(_)) => { backend.materialise(&id).await?; backend.get(&id, None).await } other => other, } } 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()); } } #[cfg(test)] mod catalog_guard_tests { //! What `sync_catalog` does when it cannot read the server's copy. //! //! The bug these exist for was a control-flow one — `if let Ok(bytes)` //! folding every failure into "there is none yet" and falling through to //! the upload — so the thing to assert is not a value but *whether a write //! happened at all*. use super::*; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; /// A backend whose read fails in a chosen way, counting writes. struct Fussy { fail_with: Option, puts: Arc, caps: dr_sync::Capabilities, } impl Fussy { fn reading(fail_with: Option) -> (Self, Arc) { let puts = Arc::new(AtomicUsize::new(0)); ( Self { fail_with, puts: puts.clone(), caps: dr_sync::Capabilities::minimal(), }, puts, ) } } #[async_trait::async_trait] impl RemoteBackend for Fussy { fn capabilities(&self) -> &dr_sync::Capabilities { &self.caps } fn name(&self) -> &str { "fussy" } async fn list( &self, _dir: &RemotePath, _since: Option<&dr_sync::Validator>, ) -> Result, RemoteError> { Ok(Vec::new()) } async fn dir_validator( &self, _dir: &RemotePath, ) -> Result { Err(RemoteError::Unsupported("test")) } async fn delta( &self, _c: &dr_sync::Cursor, ) -> Result<(Vec, dr_sync::Cursor), RemoteError> { Err(RemoteError::Unsupported("test")) } async fn get( &self, _id: &RemoteId, _r: Option>, ) -> Result, RemoteError> { match &self.fail_with { Some(RemoteError::NotFound(s)) => Err(RemoteError::NotFound(s.clone())), Some(RemoteError::NotMaterialised(s)) => { Err(RemoteError::NotMaterialised(s.clone())) } Some(_) => Err(RemoteError::PermissionDenied), None => Ok(Vec::new()), } } async fn put( &self, _p: &RemotePath, _b: Vec, _pc: Option, ) -> Result { self.puts.fetch_add(1, Ordering::SeqCst); Ok(dr_sync::Validator::new("v")) } async fn delete( &self, _id: &RemoteId, _pc: Option, ) -> Result<(), RemoteError> { Ok(()) } async fn move_to(&self, _f: &RemoteId, _t: &RemotePath) -> Result<(), RemoteError> { Ok(()) } async fn create_dir(&self, _p: &RemotePath) -> Result<(), RemoteError> { Ok(()) } } /// A real catalog and a scratch directory, since `sync_catalog` snapshots /// one before uploading. fn fixture(name: &str) -> (std::path::PathBuf, std::path::PathBuf) { let dir = std::env::temp_dir().join(format!("dr-catalog-guard-{name}")); let _ = std::fs::remove_dir_all(&dir); std::fs::create_dir_all(dir.join("scratch")).unwrap(); let catalog_path = dir.join("catalog.sqlite"); dr_catalog::Catalog::open(&catalog_path).unwrap(); (catalog_path, dir.join("scratch")) } async fn run_with(fail_with: Option, name: &str) -> (usize, SyncReport) { let (catalog_path, scratch) = fixture(name); let (backend, puts) = Fussy::reading(fail_with); let mut report = SyncReport::default(); sync_catalog( &backend, &RemotePath::new(".darkroom-derived"), &catalog_path, &scratch, &mut report, ) .await .unwrap(); let _ = std::fs::remove_dir_all(catalog_path.parent().unwrap()); (puts.load(Ordering::SeqCst), report) } #[tokio::test] async fn a_catalog_that_is_here_but_not_downloaded_is_never_written_over() { // The bug. On a placeholder library the snapshot is dehydrated, the // read fails, and the old code took that for "there is no remote // catalog" and pushed ours — discarding the other device's // collections and their members on every single sync. let (puts, report) = run_with( Some(RemoteError::NotMaterialised("catalog.sqlite".into())), "notmaterialised", ) .await; assert_eq!(puts, 0, "must not upload over a catalog it could not read"); assert!(!report.catalog_uploaded); assert!(!report.catalog_merged); } #[tokio::test] async fn a_catalog_that_cannot_be_read_at_all_is_never_written_over() { // Not only placeholders: a refused read, a dropped connection. Any // failure that is not "there is none" leaves the server's copy alone. let (puts, _) = run_with(Some(RemoteError::PermissionDenied), "denied").await; assert_eq!(puts, 0); } #[tokio::test] async fn the_first_sync_of_a_library_still_uploads() { // The other half, and the reason `NotFound` had to stay distinct: with // genuinely nothing on the server, ours *is* the whole truth and // refusing to push it would mean the catalog never syncs at all. let (puts, report) = run_with(Some(RemoteError::NotFound("nope".into())), "firstrun").await; assert_eq!(puts, 1, "nothing to merge, so ours goes up"); assert!(report.catalog_uploaded); } }