//! Writing local edits out to the catalog's sidecar outbox, and draining //! that outbox to the remote once a connection is available. use crate::sidecar_cache::SidecarCache; use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath}; use std::path::PathBuf; use std::sync::mpsc::Receiver; use super::scan::now_secs; /// 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, a pick/reject flag and a colour label — the cull. Judgement { rating: u8, flag: u8, label: 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, /// TRACES: FR-DEV-3f /// The film stock, when this is an image's own edit being written back. /// /// Two levels of `Option`, and both are load-bearing. The outer says /// whether this write concerns the film at all — a paste does not, /// exactly as it carries no masks. The inner is the choice itself, and /// `Some(None)` is a real edit: "develop this normally again". Without /// the distinction, clearing a film could never be saved. film: Option>, /// TRACES: FR-DEV-5 /// The named snapshots, when this is an image's own edit being /// written back: the ones the session holds, and the ids it deleted. /// A paste carries none — a snapshot is a state of one photograph. snapshots: Option<(Vec, Vec)>, /// TRACES: FR-DEV-3 | FR-CAT-8 /// The local adjustments, when this is an image's own edit being /// written back rather than a paste onto someone else's. /// /// A `Preset` is a parameter map, and a mask is not a parameter — it /// is a rule about *where*, with a chain of its own. So a save that /// carried only the preset wrote the sliders and silently dropped /// every local adjustment: the sidecar format has stored masks since /// they were added and `Version::apply` restores them, but nothing /// ever put any there. The mask survived until the session ended and /// then did not exist. /// /// `None` for a paste, which must not carry the source image's masks /// onto the target: a mask is drawn against one photograph and means /// nothing on another, and `Scope` cannot express that because it /// filters parameters. masks: Option, }, } /// 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( conn: Connection, writes: Vec, cache_dir: PathBuf, offline: bool, ) -> Receiver { 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) => { // Warn, not debug. This is unsynced user work — a rating or an // edit that exists only on this device — and the path is // the only thing that says *which* photograph and *where* // the server refused it. Filtered out at the default // level, a 403 on one file is indistinguishable from a // whole library failing. log::warn!("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 crate::remote::connect(&conn) { 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) => { // Warn, not debug. This is unsynced user work — a rating or an // edit that exists only on this device — and the path is // the only thing that says *which* photograph and *where* // the server refused it. Filtered out at the default // level, a 403 on one file is indistinguishable from a // whole library failing. log::warn!("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)] pub(super) 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)] pub(super) struct SidecarReport { written: usize, queued: usize, failed: usize, last_error: Option, } /// 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, }, } /// 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. pub(super) fn amend(base: dr_pipeline::Sidecar, w: &SidecarWrite) -> dr_pipeline::Sidecar { let mut sidecar = base; // TRACES: FR-NC-8 | FR-NC-9 // Close any split this file already carries, *before* looking for our own // version, and onto the uuid this write is about to use. // // Every device used to mint its own uuid for the same photograph, so a // frame edited on two of them holds two `default = 1` blocks and the // lookup below misses both — adding a third rather than amending either. // Fusing first folds them into one under `w.version_uuid`, which turns the // miss into a hit and makes this an amendment of the other device's work // instead of a rival to it. // // Idempotent: a file with one default and the right uuid is returned // byte-identical, so this costs nothing on the ordinary write. sidecar.fuse_default_versions(Some(&w.version_uuid)); // Amend the version this write belongs to, creating it if the file did // not have one — a photograph nobody has edited anywhere. 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, label, } => { version.rating = *rating; version.flag = *flag; version.label = *label; } Amendment::Settings { preset, scope, masks, film, .. } => { preset.amend(&mut version.params, *scope); // TRACES: FR-DEV-3f // Wholesale, like the masks below and for the same reason: this is // the whole of the image's own choice as it stands, so clearing a // film has to leave the sidecar too. if let Some(film) = film { version.film = film.clone(); } // Replaced wholesale rather than merged: this is the whole of the // image's local adjustment stack as it stands, so a layer the user // deleted has to leave the sidecar too. Cross-device merging of // two stacks is `Sidecar::merge`'s job and happens on sync, not // here (FR-NC-9). if let Some(masks) = masks { version.masks = masks.clone(); } } } // 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); // TRACES: FR-DEV-5 | FR-NC-9 // After the version, and against the uuid the fuse settled on: a // snapshot points at its edit by uuid, and the edit may have just been // renamed onto the canonical one. if let Amendment::Settings { snapshots: Some((kept, removed)), .. } = &w.amendment { sidecar.replace_snapshots(&w.version_uuid, kept.clone(), removed); } 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. pub(super) fn write_one_sidecar(cache: &SidecarCache, w: &SidecarWrite) -> Result { 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. pub(super) async fn write_one_sidecar_online( backend: &dyn RemoteBackend, cache: &SidecarCache, w: &SidecarWrite, ) -> Result { 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(conn: Connection, cache_dir: PathBuf) -> Receiver { 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 crate::remote::connect(&conn) { 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) => { // Warn, for the reason the write path does: this is // unsynced work and the path is what makes the failure // actionable. log::warn!("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. pub(super) async fn drain_one( backend: &dyn RemoteBackend, 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()); // TRACES: FR-NC-6c // A miss and a placeholder are not the same answer, and conflating them // destroys work. This read decides whether the sidecar already on the // remote is merged in; treating "the content is not on this device" as // "there is no sidecar" writes a fresh document over an existing one and // discards every edit another device put there — the exact loss the // format's unknown-key preservation exists to prevent. // // A sidecar is a few kilobytes, so the right response to a placeholder is // to fetch it, not to give up. Where that is impossible — no client // running — the entry stays queued, which is what the outbox is for. let remote = match backend.get(&id, None).await { Ok(bytes) => Some(bytes), Err(RemoteError::NotFound(_)) => None, Err(RemoteError::NotMaterialised(_)) => { backend .materialise(&id) .await .map_err(|e| format!("sidecar is not on this device ({e})"))?; match backend.get(&id, None).await { Ok(bytes) => Some(bytes), Err(e) => return Err(format!("sidecar could not be read ({e})")), } } // Anything else — a refused read, a dead connection — leaves the entry // queued rather than resolved by overwriting. Err(e) => return Err(format!("sidecar could not be read ({e})")), }; 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})")), } } } // TRACES: FR-NC-8 | FR-NC-9 // `merge_into` reconciles version by version *by uuid*, so two devices' // independently minted defaults pass straight through it and both land in // what is about to be uploaded. Fusing here is what stops the outbox from // publishing the split rather than resolving it. // // No canonical uuid: this entry may have been queued by a build that had // not derived one yet, and the smallest uuid is device-independent, which // is all convergence needs. The next write from either device moves it // onto the derived identity. local.fuse_default_versions(None); 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. pub(super) 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()), } } } #[cfg(test)] mod tests { use super::*; fn judgement(uuid: &str, rating: u8) -> SidecarWrite { SidecarWrite { image_path: "PhotosRaw/incoming/a.CR2".to_string(), version_uuid: uuid.to_string(), amendment: Amendment::Judgement { rating, flag: 0, label: 0, }, } } fn version(uuid: &str, revision: u64, modified: i64, rating: u8) -> dr_pipeline::Version { dr_pipeline::sidecar::Version { uuid: uuid.to_string(), name: "Default".to_string(), is_default: true, revision, modified, rating, ..Default::default() } } /// The regression. A sidecar already holding two independently minted /// defaults used to gain a *third* on the next write, because the lookup /// is by uuid and neither of the two was ours. #[test] fn a_write_onto_a_split_sidecar_does_not_add_a_third_version() { let mut base = dr_pipeline::Sidecar::new(); base.put(version("tablet-uuid", 1, 100, 4)); base.put(version("laptop-uuid", 2, 200, 1)); let out = amend(base, &judgement("derived-uuid", 5)); assert_eq!(out.versions.len(), 1, "the split must be closed, not grown"); assert!(out.versions.contains_key("derived-uuid")); assert_eq!(out.versions["derived-uuid"].rating, 5); } /// The other device's work has to survive the fold, or closing the split /// would be the same data loss by a different route. #[test] fn the_other_devices_edit_survives_the_write() { let mut theirs = version("tablet-uuid", 3, 300, 4); theirs .params .insert(("exposure".into(), "exposure".into()), 0.75); let mut base = dr_pipeline::Sidecar::new(); base.put(theirs); // Ours is a rating, which touches no parameter at all. let out = amend(base, &judgement("derived-uuid", 2)); let v = &out.versions["derived-uuid"]; assert_eq!( v.params.get(&("exposure".into(), "exposure".into())), Some(&0.75), "the tablet's exposure was dropped by our rating" ); assert_eq!(v.rating, 2, "and our own judgement did not land"); } /// A sidecar this device has already written must not be disturbed: the /// ordinary case is one default under the right uuid, and fusing it has to /// be a no-op beyond the amendment itself. #[test] fn the_ordinary_write_is_unaffected() { let mut base = dr_pipeline::Sidecar::new(); base.put(version("derived-uuid", 7, 700, 3)); let out = amend(base, &judgement("derived-uuid", 5)); assert_eq!(out.versions.len(), 1); assert_eq!(out.versions["derived-uuid"].rating, 5); assert_eq!( out.versions["derived-uuid"].revision, 8, "one bump for one edit" ); } /// A virtual copy is not a rival default and must be left where it is. #[test] fn a_named_version_is_not_folded_into_the_default() { let mut copy = version("for-print", 4, 400, 5); copy.is_default = false; copy.name = "For print".to_string(); let mut base = dr_pipeline::Sidecar::new(); base.put(version("tablet-uuid", 1, 100, 4)); base.put(copy); let out = amend(base, &judgement("derived-uuid", 2)); assert_eq!(out.versions.len(), 2); assert_eq!(out.versions["for-print"].name, "For print"); assert_eq!(out.versions["for-print"].rating, 5); } }