Files
DarkRoom/ui/dr-ui/src/derived_sync.rs
T
dtourolle 83f0461ce7 Sync develop presets through the library
Presets were the one piece of the photographer's work that never left
the device: faces, sidecars, albums, collections, keywords and camera
profiles all travel with the sync pass, the preset library did not.

It now goes to <derived>/presets/library.drpl. PresetLibrary::merge
decides each name against the base the last exchange left (kept per
library beside place.json), so presets added on two devices both
survive, a deletion reaches the other device instead of being restored
by it, and an edit outlives a deletion made elsewhere. The upload is
If-Match / If-None-Match on the server's copy, and a 412 reads and
merges again, so two devices exchanging at once cannot save over each
other. A server copy that will not parse (a newer build's) is left
alone, and a local file that will not read stops the exchange rather
than being taken for an empty library.

The develop view's save merges with the file when the sync changed it
since the view read it, and a sync that brought presets reloads and
redraws the list.

Also corrects the register, which still said camera profiles do not
sync.
2026-10-04 00:50:09 -04:00

2304 lines
92 KiB
Rust

//! 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 crate::executors::{self, Executor};
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,
/// The copy on the server was damaged beyond reading, and ours replaced it.
///
/// Counted because it is the one outcome here that *destroys* something. A
/// sync that silently overwrote a peer's collections would be the bug this
/// whole path is written to avoid, so when it does happen the report has to
/// be able to say so rather than looking like an ordinary push.
pub catalog_replaced: 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,
/// Images dated from the remote's snapshot rather than by this device's
/// own sweep — what makes the timeline whole on a fresh device.
pub dates_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,
/// TRACES: FR-UI-8
/// Whether the exchange moved where this device thinks the photographer is.
///
/// Counted apart from everything above because it is the one thing here
/// that is not derived state: a shard or a catalog snapshot is a faster way
/// to learn what this device could have worked out for itself, and a place
/// is a fact only the other device knew.
pub place_adopted: bool,
pub place_uploaded: bool,
/// TRACES: FR-DEV-3e
/// Camera profiles exchanged with the library's `profiles` folder
/// (camera-profiles.md §13).
pub profiles_uploaded: usize,
pub profiles_downloaded: usize,
/// TRACES: FR-DEV-6
/// The develop presets: whether another device's changes reached this
/// one's library, and whether this one's reached the server.
pub presets_adopted: bool,
pub presets_uploaded: bool,
}
impl SyncReport {
pub fn did_anything(&self) -> bool {
self.shards_uploaded > 0
|| self.shards_downloaded > 0
|| self.catalog_uploaded
|| self.catalog_merged
|| self.catalog_replaced
|| self.face_shards_uploaded > 0
|| self.face_shards_downloaded > 0
|| self.place_adopted
|| self.profiles_downloaded > 0
|| self.presets_adopted
}
}
/// 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).
// One argument per thing the pass touches; see `run` below.
#[allow(clippy::too_many_arguments)]
pub fn spawn_sync(
conn: Connection,
root: String,
thumbs_dir: PathBuf,
catalog_path: PathBuf,
place_path: PathBuf,
presets: PresetFiles,
scratch: PathBuf,
// TRACES: FR-CULL-8
// Which face pipeline's shards to export and adopt. From the settings
// page by way of the library controller, because a shard written under
// one detector's faces must not be read back as another's.
face_model_id: String,
) -> std::sync::mpsc::Receiver<SyncMessage> {
let (tx, rx) = std::sync::mpsc::channel();
executors::spawn(Executor::Network, "sync", 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,
&place_path,
&presets,
&scratch,
&face_model_id,
&tx,
)
.await
{
Ok(report) => {
let _ = tx.send(SyncMessage::Finished(Box::new(report)));
}
Err(e) => {
let _ = tx.send(SyncMessage::Failed(e));
}
}
});
});
rx
}
// Nine, because a sync touches nine distinct things — the same reason
// `repairs::spawn` carries the allow: bundling them into a struct would name
// nothing that exists.
#[allow(clippy::too_many_arguments)]
async fn run(
backend: &dyn RemoteBackend,
root: &str,
thumbs_dir: &Path,
catalog_path: &Path,
place_path: &Path,
presets: &PresetFiles,
scratch: &Path,
face_model_id: &str,
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;
// Small before large, and what the user is waiting for before what
// fills in behind them. Faces first because the catalog merge assigns
// identities to faces this device holds, so they must be here by then;
// the catalog next for collections, people and dates; thumbnails last,
// because a fresh device's thumbnail stage is hundreds of megabytes and
// everything queued behind it — the sidebar, the timeline, the names —
// was invisible for as long as it ran.
let _ = tx.send(SyncMessage::Status("checking faces…".into()));
sync_face_shards(
backend,
&base,
catalog_path,
scratch,
face_model_id,
&mut report,
tx,
)
.await?;
let _ = tx.send(SyncMessage::Status("checking collections…".into()));
sync_catalog(backend, &base, catalog_path, scratch, &mut report).await?;
let _ = tx.send(SyncMessage::Status("checking thumbnails…".into()));
sync_shards(backend, &base, thumbs_dir, scratch, &mut report).await?;
// TRACES: FR-DEV-3e
// The camera profiles, so a profile copied out of a DNG on one device
// renders that body's raws on every device (camera-profiles.md §13).
// After the catalog and before the place; like the place it never fails
// the pass.
if let Some(dir) = dr_decode::dcp::profiles_directory() {
let _ = tx.send(SyncMessage::Status("checking camera profiles…".into()));
sync_profiles(backend, &base, &dir, &mut report).await;
}
// TRACES: FR-DEV-6
// The develop presets, a few kilobytes. Like the profiles, never fails
// the pass.
let _ = tx.send(SyncMessage::Status("checking presets…".into()));
sync_presets(backend, &base, presets, &mut report).await;
// TRACES: FR-UI-8
// Last, and it costs one small GET plus at most one small PUT. Last because
// it is the only thing here that is not derived state and so the only thing
// whose loss the user would not notice: a pass that ran out of connectivity
// should spend what it had on the shards and the catalog.
sync_place(backend, &base, place_path, &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<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 {
// A snapshot, never the live file — see `ThumbStore::snapshot_shard`
// for the zero-byte upload that reading the file produced.
let snapshot = scratch.join(format!("thumb-shard-{:04}-upload.sqlite", shard.id));
let bytes = match store.snapshot_shard(shard.id, &snapshot) {
Ok(()) => std::fs::read(&snapshot),
Err(e) => Err(std::io::Error::other(e.to_string())),
};
let _ = std::fs::remove_file(&snapshot);
let bytes = match bytes {
Ok(b) => b,
Err(e) => {
log::warn!("snapshotting thumbnail shard {}: {e}", shard.id);
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,
face_model_id: &str,
report: &mut SyncReport,
tx: &std::sync::mpsc::Sender<SyncMessage>,
) -> 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 = face_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<String, u64> = 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 and its people.
///
/// Only collections, keywords and people 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.
///
/// # Generations
///
/// TRACES: FR-NC-9 | NFR-R2
/// The server keeps the current copy and the [`GENERATIONS`] before it:
/// `catalog.sqlite`, then `catalog.1.sqlite` (the one it replaced), `.2`,
/// `.3`. A push uploads to a temporary name, rotates, and moves the upload into
/// place — so at no moment is there no current copy, and a copy that turns out
/// damaged has the one before it to fall back on. Rotation is server-side
/// renames; the only transfer is the upload itself.
///
/// This is what a damaged copy used to lack. With one copy and nothing behind
/// it, "the current file will not open" left two answers, both bad: refuse for
/// ever, or overwrite with ours and lose whatever another device had added
/// since. Now it has a third — merge from the newest readable generation, which
/// loses nothing — and the damaged file itself is kept as `.1` by the ordinary
/// rotation rather than by a separate upload.
async fn sync_catalog(
backend: &dyn RemoteBackend,
base: &RemotePath,
catalog_path: &Path,
scratch: &Path,
report: &mut SyncReport,
) -> Result<(), String> {
let remote_name = CATALOG_NAME;
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 catalog = match dr_catalog::Catalog::open(catalog_path) {
Ok(c) => c,
Err(e) => {
log::warn!("not pushing the catalog: opening ours to merge: {e}");
return Ok(());
}
};
match merge_downloaded(&catalog, scratch, &bytes, report) {
Ok(()) => {}
// TRACES: FR-NC-9
// Damaged beyond reading, and the whole file arrived — so it is
// not a short download, and no device will read it either. Fall
// back to the generation before it; only with none readable is
// ours the whole truth. See the header for why refusing here for
// ever was the wrong answer.
Err(dr_catalog::CatalogError::Corrupt { detail }) => {
if !arrived_whole(backend, base, remote_name, bytes.len(), &detail).await {
return Ok(());
}
match merge_from_generations(backend, base, &catalog, scratch, report).await {
Some(n) => log::warn!(
"the catalog on the server is damaged ({detail}) and all {} bytes of \
it arrived, so no device can read it; merged from the generation \
before it ({}) instead, and replacing it",
bytes.len(),
generation_name(n)
),
None => log::warn!(
"the catalog on the server is damaged ({detail}) and all {} bytes of \
it arrived, so no device can read it, and no earlier generation is \
readable either; replacing it with this device's copy",
bytes.len()
),
}
report.catalog_replaced = true;
}
// Unreadable is not the same as absent: it may be a newer
// format. Ours must not go over it.
Err(e) => {
log::warn!("not pushing the catalog: merging the server's copy: {e}");
return Ok(());
}
}
}
// ---- 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())?;
let _ = std::fs::remove_file(&snapshot);
// To a temporary name first. The upload is the only step that can fail
// half-way, and a half-uploaded *current* copy is exactly the damaged file
// this whole scheme exists to survive. Under its own name a failure leaves
// the current copy untouched and costs one stray file, retried next pass.
let staging = RemotePath::new(format!("{}/{UPLOAD_NAME}", base.as_str()));
let sent = bytes.len() as u64;
if let Err(e) = backend.put(&staging, bytes, None).await {
log::warn!("uploading catalog: {e}");
return Ok(());
}
// TRACES: NFR-R2
// Confirm the server holds what was sent before it becomes the copy every
// other device reads. A chunked upload is assembled server-side, and an
// assembly that went wrong is a file of plausible size that no device can
// open — the one failure the generations exist to survive, and cheaper
// to catch here, on the device that caused it, than on every other one
// after. One listing; the size is what the server can vouch for without
// reading the file back.
match remote_size(backend, base, UPLOAD_NAME).await {
Some(held) if held == sent => {}
Some(held) => {
log::warn!(
"not replacing the catalog: sent {sent} bytes but the server holds {held}; \
the upload is discarded and retried next pass"
);
let _ = backend.delete(&RemoteId::Path(staging), None).await;
return Ok(());
}
None => {
log::warn!(
"not replacing the catalog: the server would not confirm the upload's size; \
it is discarded and retried next pass"
);
let _ = backend.delete(&RemoteId::Path(staging), None).await;
return Ok(());
}
}
// Rotate, then move the upload into place. Every step here is a rename on
// the server, and every destination is empty by the time it is written to
// — a `move_to` will not overwrite, by design — so a failure at any point
// leaves a gap in the generations and never a missing current copy.
if let Err(e) = rotate_generations(backend, base).await {
log::warn!("not replacing the catalog: rotating the earlier copies: {e}");
return Ok(());
}
match backend.move_to(&RemoteId::Path(staging), &target).await {
Ok(()) => report.catalog_uploaded = true,
Err(e) => log::warn!("moving the uploaded catalog into place: {e}"),
}
Ok(())
}
/// The current copy's name on the server.
const CATALOG_NAME: &str = "catalog.sqlite";
/// Where a push lands before it is rotated into place.
const UPLOAD_NAME: &str = "catalog.upload.sqlite";
/// How many earlier copies the server keeps behind the current one.
///
/// Three, because what they are for is surviving one damaged push and the one
/// or two syncs it may take for a device to notice. More would cost nothing in
/// transfer — rotation is renames — but each is a 40 MB file on the account's
/// quota, and the local backups (NFR-R2) are the long-term store.
const GENERATIONS: usize = 3;
/// `catalog.N.sqlite` for `1 <= N <= GENERATIONS`; `.1` is the newest.
fn generation_name(n: usize) -> String {
format!("catalog.{n}.sqlite")
}
/// Merge a downloaded catalog into ours, through a file in scratch.
///
/// Attaching needs a path, and the download is bytes. Written and removed here
/// so the callers — the current copy, and each generation tried after it —
/// cannot disagree about cleanup.
fn merge_downloaded(
catalog: &dr_catalog::Catalog,
scratch: &Path,
bytes: &[u8],
report: &mut SyncReport,
) -> Result<(), dr_catalog::CatalogError> {
let downloaded = scratch.join("catalog-remote.sqlite");
std::fs::write(&downloaded, bytes).map_err(|e| dr_catalog::CatalogError::Io(e.to_string()))?;
let result = catalog.merge_remote_catalog(&downloaded);
let _ = std::fs::remove_file(&downloaded);
let merge = result?;
report.catalog_merged = true;
report.collections_gained += merge.inserted + merge.updated;
report.members_gained += merge.members_added;
report.dates_gained += merge.metadata_adopted;
Ok(())
}
/// TRACES: FR-NC-9
/// Merge from the newest generation that reads, when the current copy will
/// not. Returns which one, or `None` when none of them does.
///
/// Newest first, and the first readable one wins: a generation is a complete
/// snapshot, so an older one adds nothing a newer one lacks. A generation that
/// is absent, damaged, or from a newer schema is skipped the same way — none of
/// those is a reason to stop looking further back.
async fn merge_from_generations(
backend: &dyn RemoteBackend,
base: &RemotePath,
catalog: &dr_catalog::Catalog,
scratch: &Path,
report: &mut SyncReport,
) -> Option<usize> {
for n in 1..=GENERATIONS {
let path = RemotePath::new(format!("{}/{}", base.as_str(), generation_name(n)));
let bytes = match read_derived(backend, &path).await {
Ok(b) => b,
Err(RemoteError::NotFound(_)) => continue,
Err(e) => {
log::debug!("skipping {}: {e}", generation_name(n));
continue;
}
};
match merge_downloaded(catalog, scratch, &bytes, report) {
Ok(()) => return Some(n),
Err(e) => log::debug!("skipping {}: {e}", generation_name(n)),
}
}
None
}
/// Make room for a new current copy: drop the oldest generation and shift the
/// rest back by one, ending with the current copy as `.1`.
///
/// Oldest first, so that each destination is empty when it is moved into —
/// `move_to` refuses to overwrite, and rightly. A name that is not there is
/// not an error at any step: a library that has synced twice has no `.3` yet.
async fn rotate_generations(backend: &dyn RemoteBackend, base: &RemotePath) -> Result<(), String> {
let at = |name: String| RemotePath::new(format!("{}/{name}", base.as_str()));
match backend
.delete(&RemoteId::Path(at(generation_name(GENERATIONS))), None)
.await
{
Ok(()) | Err(RemoteError::NotFound(_)) => {}
Err(e) => return Err(format!("dropping {}: {e}", generation_name(GENERATIONS))),
}
for n in (1..GENERATIONS).rev() {
match backend
.move_to(
&RemoteId::Path(at(generation_name(n))),
&at(generation_name(n + 1)),
)
.await
{
Ok(()) | Err(RemoteError::NotFound(_)) => {}
Err(e) => return Err(format!("moving {} back: {e}", generation_name(n))),
}
}
match backend
.move_to(
&RemoteId::Path(at(CATALOG_NAME.into())),
&at(generation_name(1)),
)
.await
{
Ok(()) | Err(RemoteError::NotFound(_)) => Ok(()),
Err(e) => Err(format!("setting the current copy back: {e}")),
}
}
/// TRACES: FR-UI-8
/// The place's name inside the derived folder.
///
/// The same name the local copy has, so `.darkroom-derived/place.json` and the
/// file beside the catalog are visibly the same thing. There is one of these per
/// library, not one per device: the question it answers — where is the
/// photographer — has one answer, and a folder of per-device files would need a
/// listing and N fetches to work out which of them was current.
const PLACE_NAME: &str = "place.json";
/// TRACES: FR-UI-8
/// Exchange the place with the server. The newer record wins, in both
/// directions.
///
/// # Why this cannot clobber the way the catalog could
///
/// [`sync_catalog`] has to refuse to upload when it cannot read the server's
/// copy, because its upload is a read-modify-write: writing without merging
/// discards the other device's collections. This is not that. A place is
/// *replaced*, never merged, so there is nothing of theirs inside ours to lose.
///
/// It still gives up on an unreadable read rather than uploading over it, for a
/// smaller reason: a server that will not answer a GET is not one to spend a PUT
/// on, and a record we could not compare against might be newer than ours —
/// overwriting it would move the other device's photographer without ever having
/// seen where they were.
///
/// # Failures are not propagated
///
/// This returns nothing and takes no `?`. Every other step in [`run`] carries
/// state that has to arrive; this one carries a scroll position, and a sync that
/// reported itself failed — putting an error in front of the user and skipping
/// nothing, since it runs last — because a position file could not be written
/// would be reporting the wrong thing entirely.
async fn sync_place(
backend: &dyn RemoteBackend,
base: &RemotePath,
place_path: &Path,
report: &mut SyncReport,
) {
let target = RemotePath::new(format!("{}/{PLACE_NAME}", base.as_str()));
let theirs = match read_derived(backend, &target).await {
Ok(bytes) => crate::place::from_bytes(&bytes, "the place on the server"),
// Nobody has recorded one for this library yet.
Err(RemoteError::NotFound(_)) => None,
Err(e) => {
log::debug!("not exchanging the place: reading the server's copy: {e}");
return;
}
};
let store = crate::place::PlaceStore::open_at(place_path.to_path_buf());
// Take theirs first, so what goes back up is whichever of the two is
// current rather than always ours.
if let Some(theirs) = theirs.as_ref() {
if store.adopt(theirs) {
log::info!(
"the place moved to where {} left off",
if theirs.device.is_empty() {
"another device"
} else {
&theirs.device
}
);
report.place_adopted = true;
}
}
// Then push, if there is anything to say. Re-read rather than reusing what
// was loaded above: `adopt` may have just replaced it, and uploading a
// record the server already has is a PUT for nothing.
let Some(ours) = store.load() else {
return;
};
if !ours.supersedes(theirs.as_ref()) {
return;
}
let bytes = match crate::place::to_bytes(&ours) {
Ok(b) => b,
Err(e) => {
log::debug!("serialising the place: {e}");
return;
}
};
match backend.put(&target, bytes, None).await {
Ok(_) => report.place_uploaded = true,
Err(e) => log::debug!("uploading the place: {e}"),
}
}
/// TRACES: FR-DEV-3e
/// Exchange camera profiles with `<derived>/profiles` (camera-profiles.md
/// §13).
///
/// Profiles are immutable and named for what they hold (`<camera>
/// <profile>.dcp`, `dcp::save`), so a name and a size say whether two copies
/// are the same: upload what the server lacks or holds at another size,
/// download what this device lacks. A download is parsed before it is kept,
/// and written beside its name then renamed, because the profile loader
/// reads every `.dcp` in the folder and a half-written one would be skipped
/// with a warning rather than retried. Never fails the pass: a profile that
/// did not travel this time travels next time.
async fn sync_profiles(
backend: &dyn RemoteBackend,
base: &RemotePath,
local_dir: &Path,
report: &mut SyncReport,
) {
let dir = RemotePath::new(format!("{}/profiles", base.as_str()));
let is_profile = |name: &str| {
Path::new(name)
.extension()
.is_some_and(|x| x.eq_ignore_ascii_case("dcp"))
};
let remote: std::collections::HashMap<String, u64> = backend
.list(&dir, None)
.await
.map(|entries| {
entries
.into_iter()
.filter(|e| e.kind == dr_sync::EntryKind::File && is_profile(e.path.name()))
.map(|e| (e.path.name().to_string(), e.size))
.collect()
})
// Absent until the first device uploads one; an empty answer.
.unwrap_or_default();
let local: std::collections::HashMap<String, PathBuf> = std::fs::read_dir(local_dir)
.map(|entries| {
entries
.flatten()
.filter_map(|e| {
let name = e.file_name().to_string_lossy().into_owned();
is_profile(&name).then(|| (name, e.path()))
})
.collect()
})
.unwrap_or_default();
// ---- upload ----------------------------------------------------------
let to_upload: Vec<(&String, &PathBuf)> = local
.iter()
.filter(|(name, path)| {
let size = std::fs::metadata(path).map(|m| m.len()).unwrap_or(0);
remote.get(*name) != Some(&size)
})
.collect();
if !to_upload.is_empty() {
let _ = backend.create_dir(&dir).await;
}
for (name, path) in to_upload {
let Ok(bytes) = std::fs::read(path) else {
continue;
};
let target = RemotePath::new(format!("{}/{name}", dir.as_str()));
match backend.put(&target, bytes, None).await {
Ok(_) => report.profiles_uploaded += 1,
Err(e) => log::debug!("uploading camera profile {name}: {e}"),
}
}
// ---- download --------------------------------------------------------
for name in remote.keys().filter(|n| !local.contains_key(*n)) {
let source = RemotePath::new(format!("{}/{name}", dir.as_str()));
let bytes = match read_derived(backend, &source).await {
Ok(b) => b,
Err(e) => {
log::debug!("fetching camera profile {name}: {e}");
continue;
}
};
if let Err(why) = dr_decode::dcp::Dcp::parse(&bytes) {
log::warn!("camera profile {name} on the server is not one: {why}");
continue;
}
let path = local_dir.join(name);
let partial = local_dir.join(format!("{name}.part"));
let written = std::fs::create_dir_all(local_dir)
.and_then(|()| std::fs::write(&partial, &bytes))
.and_then(|()| std::fs::rename(&partial, &path));
match written {
Ok(()) => report.profiles_downloaded += 1,
Err(e) => {
let _ = std::fs::remove_file(&partial);
log::warn!("saving camera profile {name}: {e}");
}
}
}
if report.profiles_downloaded > 0 {
log::info!(
"camera profiles: {} fetched from the library",
report.profiles_downloaded
);
// Decodes from now on see them; one already in flight keeps the set
// it started with.
dr_decode::dcp::set_profiles_directory(local_dir.to_path_buf());
}
}
/// TRACES: FR-DEV-6
/// Where this device keeps the preset library, and what the last exchange
/// of it with this library left both sides holding.
///
/// The library is the device's, shared by every library it opens; the base
/// is per library, because each library's server holds its own copy and has
/// its own history of exchanges with this device.
#[derive(Debug, Clone)]
pub struct PresetFiles {
pub library: PathBuf,
pub base: PathBuf,
}
/// The preset library on the server, in its own folder so finding it is a
/// listing of one file rather than of every shard beside it.
const PRESETS_DIR: &str = "presets";
const PRESETS_NAME: &str = "library.drpl";
/// TRACES: FR-DEV-6
/// Exchange the develop presets with `<derived>/presets/library.drpl`.
///
/// Unlike the place, this is merged rather than replaced: a preset saved on
/// the tablet and another saved on the desktop between two passes must both
/// survive, and a preset deleted on one must not come back from the other.
/// [`PresetLibrary::merge`](dr_pipeline::PresetLibrary::merge) decides each
/// name against the base the last exchange left, which is what tells a
/// deletion here from an addition there.
///
/// The write is conditional on the server still holding what was read, so
/// two devices exchanging at once cannot each save over the other's
/// additions; the one that loses the race reads again and merges again.
/// Never fails the pass.
async fn sync_presets(
backend: &dyn RemoteBackend,
base: &RemotePath,
files: &PresetFiles,
report: &mut SyncReport,
) {
for _ in 0..3 {
match exchange_presets(backend, base, files, report).await {
Err(RemoteError::PreconditionFailed) => {
log::debug!("presets: another device wrote first; reading again");
}
Err(e) => {
log::debug!("not exchanging presets: {e}");
return;
}
Ok(()) => return,
}
}
}
async fn exchange_presets(
backend: &dyn RemoteBackend,
base: &RemotePath,
files: &PresetFiles,
report: &mut SyncReport,
) -> Result<(), RemoteError> {
use crate::preset_store::PresetStore;
use dr_pipeline::PresetLibrary;
let dir = RemotePath::new(format!("{}/{PRESETS_DIR}", base.as_str()));
let target = RemotePath::new(format!("{}/{PRESETS_NAME}", dir.as_str()));
// A listing that fails reads as "not there". That is safe because the
// write below is then `IfAbsent`, which a server holding one refuses.
let listed = backend.list(&dir, None).await.ok().and_then(|entries| {
entries
.into_iter()
.find(|e| e.kind == dr_sync::EntryKind::File && e.path.name() == PRESETS_NAME)
});
let mut theirs = match &listed {
None => PresetLibrary::default(),
Some(_) => {
let bytes = read_derived(backend, &target).await?;
match std::str::from_utf8(&bytes)
.map_err(|e| e.to_string())
.and_then(|t| PresetLibrary::parse(t).map_err(|e| e.to_string()))
{
Ok(library) => library,
Err(e) => {
// A newer build's format, or damage. Either way not ours
// to write over.
log::warn!("the preset library on the server will not read ({e}); left alone");
return Ok(());
}
}
}
};
let store = PresetStore::open_at(files.library.clone());
let read = match store.try_load() {
Ok(library) => library.unwrap_or_default(),
Err(e) => {
// An empty library here would read as every preset deleted.
log::warn!("not exchanging presets: {}: {e}", store.path().display());
return Ok(());
}
};
let base_store = PresetStore::open_at(files.base.clone());
let last = base_store.try_load().ok().flatten().unwrap_or_default();
// Copies an older build seeded of the shipped presets are not the
// photographer's, and are dropped on both sides before anything travels.
let mut ours = read.clone();
dr_pipeline::bundled::forget_unchanged_copies(&mut ours);
dr_pipeline::bundled::forget_unchanged_copies(&mut theirs);
let merged = PresetLibrary::merge(&last, &ours, &theirs);
if listed.is_none() && merged.is_empty() {
return Ok(());
}
if listed.is_none() || merged != theirs {
let _ = backend.create_dir(&dir).await;
let precondition = match &listed {
Some(entry) => dr_sync::Precondition::IfMatch(entry.validator.clone()),
None => dr_sync::Precondition::IfAbsent,
};
backend
.put(&target, merged.to_text().into_bytes(), Some(precondition))
.await?;
report.presets_uploaded = true;
}
if merged != ours {
// The develop view saves this file too. A preset it saved since the
// read above is merged in rather than written over; it reaches the
// server on the next pass, as an addition against the base below.
let local = match store.try_load() {
Ok(Some(now)) if now != read => PresetLibrary::merge(&read, &merged, &now),
_ => merged.clone(),
};
match store.save(&local) {
Ok(()) => report.presets_adopted = true,
Err(e) => {
log::warn!("saving presets from the library: {e}");
return Ok(());
}
}
}
if merged != last {
if let Err(e) = base_store.save(&merged) {
log::debug!("recording the preset exchange: {e}");
}
}
Ok(())
}
/// TRACES: FR-UI-8
/// Fetch just the place, for the handover at launch.
///
/// # Why this is not simply [`sync_place`]
///
/// The full pass runs after a thumbnail sweep or when the Sync button is
/// pressed, neither of which happens on an ordinary launch — so a place left on
/// the tablet would reach the desktop one launch late, which is one launch too
/// many for a feature whose whole claim is that you pick up where you stopped.
///
/// This is the small half: one GET of a few hundred bytes, started beside the
/// scan rather than after it, so the answer is usually in hand before the user
/// has decided what to look at. It writes nothing and uploads nothing — the
/// exchange proper still happens in [`sync_place`], which is where "ours is
/// newer, push it" belongs.
///
/// `Ok(None)` is "there is not one there", which on a first launch against a
/// fresh library is the normal answer and not worth a word to the user.
pub fn spawn_place_fetch(
conn: Connection,
root: String,
) -> std::sync::mpsc::Receiver<Result<Option<dr_types::Place>, String>> {
let (tx, rx) = std::sync::mpsc::channel();
executors::spawn(Executor::Network, "places", move || {
let rt = match crate::net_runtime::build() {
Ok(rt) => rt,
Err(e) => {
let _ = tx.send(Err(e.to_string()));
return;
}
};
rt.block_on(async {
let backend = match crate::remote::connect(&conn) {
Ok(b) => b,
Err(e) => {
let _ = tx.send(Err(e.to_string()));
return;
}
};
let base = derived_path(&root);
let target = RemotePath::new(format!("{}/{PLACE_NAME}", base.as_str()));
let got = match read_derived(&*backend, &target).await {
Ok(bytes) => Ok(crate::place::from_bytes(&bytes, "the place on the server")),
Err(RemoteError::NotFound(_)) => Ok(None),
Err(e) => Err(e.to_string()),
};
let _ = tx.send(got);
});
});
rx
}
/// 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<Vec<u8>, 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,
}
}
/// TRACES: FR-NC-9
/// Whether a download that will not open is the server's whole file.
///
/// A truncated download is also unreadable, and on a phone it is the far
/// likelier story — this library's logs are full of aborted bodies and DNS
/// failures. Treating that as damage would let one bad connection discard a
/// catalog the server was holding perfectly well. So the size the server
/// advertises is compared against what actually arrived, and anything short,
/// or any size the listing cannot confirm, is the transport failure it is:
/// the caller must leave the server's copy alone.
async fn arrived_whole(
backend: &dyn RemoteBackend,
base: &RemotePath,
name: &str,
received: usize,
detail: &str,
) -> bool {
let Some(advertised) = remote_size(backend, base, name).await else {
log::warn!(
"not pushing the catalog: the copy on the server will not open ({detail}), \
but its size could not be confirmed, so it may simply have arrived short"
);
return false;
};
if advertised != received as u64 {
log::warn!(
"not pushing the catalog: the copy on the server will not open ({detail}), \
but only {received} of {advertised} bytes arrived — that is a truncated \
download, not a damaged file, so the server's copy is left alone"
);
return false;
}
true
}
/// The size the server says an entry in the derived folder has.
///
/// One listing of a folder that holds a handful of files, rather than a HEAD
/// the backend trait does not offer. `None` covers "the listing failed", "it is
/// not there", and "the size it reports means nothing" alike, and the caller
/// treats all of them as not knowing — which is the answer that declines to
/// overwrite.
///
/// A placeholder is excluded rather than trusted: `RemoteEntry::size` is
/// explicitly not meaningful when `materialised` is false — a suffix-mode stub
/// is one byte and carries no record of what it stands for — so comparing a
/// download against it would be comparing against nothing.
async fn remote_size(backend: &dyn RemoteBackend, base: &RemotePath, name: &str) -> Option<u64> {
backend
.list(base, None)
.await
.ok()?
.into_iter()
.find(|e| e.kind == dr_sync::EntryKind::File && e.path.name() == name)
.filter(|e| e.materialised)
.map(|e| e.size)
}
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 derived_guard_tests {
//! What the two read-then-write steps do when they cannot read the server's
//! copy — the catalog snapshot, and the place.
//!
//! 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<RemoteError>,
/// What a successful read returns. Empty for the catalog tests, which
/// only ever exercise failures; the place tests need real content,
/// because "theirs is newer than ours" is the branch under test.
body: Vec<u8>,
puts: Arc<AtomicUsize>,
/// What the last successful write carried, so a test can assert *which*
/// record went up rather than only that one did.
last_put: Arc<std::sync::Mutex<Vec<u8>>>,
caps: dr_sync::Capabilities,
/// The size `list` claims `catalog.sqlite` has, if it lists it at all.
/// This is what says whether a body that will not open is a damaged
/// file or merely a short download.
advertise: Option<u64>,
/// Bodies by file name, for tests that need the current copy and a
/// generation to differ. When empty every read serves `body`; when
/// not, a name absent here reads as `NotFound`.
bodies: std::collections::HashMap<String, Vec<u8>>,
/// Every `move_to`, as (from, to) names, in order.
moves: Arc<std::sync::Mutex<Vec<(String, String)>>>,
/// What `list` claims the staged upload's size is, when a test wants
/// the server to have assembled it wrongly. `None` reports the truth.
staged_size: Option<u64>,
}
impl Fussy {
fn reading(fail_with: Option<RemoteError>) -> (Self, Arc<AtomicUsize>) {
let (f, puts, _) = Self::serving(fail_with, Vec::new());
(f, puts)
}
fn serving(
fail_with: Option<RemoteError>,
body: Vec<u8>,
) -> (Self, Arc<AtomicUsize>, Arc<std::sync::Mutex<Vec<u8>>>) {
let puts = Arc::new(AtomicUsize::new(0));
let last_put = Arc::new(std::sync::Mutex::new(Vec::new()));
(
Self {
fail_with,
body,
puts: puts.clone(),
last_put: last_put.clone(),
caps: dr_sync::Capabilities::minimal(),
advertise: None,
bodies: Default::default(),
moves: Default::default(),
staged_size: None,
},
puts,
last_put,
)
}
}
#[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<Vec<dr_sync::RemoteEntry>, RemoteError> {
let entry = |name: &str, size: u64| {
let path = RemotePath::new(format!("{}/{name}", dir.as_str()));
dr_sync::RemoteEntry {
id: RemoteId::Path(path.clone()),
path,
kind: dr_sync::EntryKind::File,
validator: dr_sync::Validator::new("v"),
size,
modified: None,
has_preview: false,
materialised: true,
}
};
let mut out = Vec::new();
if let Some(size) = self.advertise {
out.push(entry(CATALOG_NAME, size));
}
// The staged upload lists at the size of what was last put, as a
// server that assembled it correctly would report — unless a test
// says the assembly went wrong.
if self.puts.load(Ordering::SeqCst) > 0 {
let held = self
.staged_size
.unwrap_or(self.last_put.lock().unwrap().len() as u64);
out.push(entry(UPLOAD_NAME, held));
}
Ok(out)
}
async fn dir_validator(
&self,
_dir: &RemotePath,
) -> Result<dr_sync::Validator, RemoteError> {
Err(RemoteError::Unsupported("test"))
}
async fn delta(
&self,
_c: &dr_sync::Cursor,
) -> Result<(Vec<dr_sync::RemoteChange>, dr_sync::Cursor), RemoteError> {
Err(RemoteError::Unsupported("test"))
}
async fn get(
&self,
id: &RemoteId,
_r: Option<std::ops::Range<u64>>,
) -> Result<Vec<u8>, 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 if self.bodies.is_empty() => Ok(self.body.clone()),
None => {
let name = match id {
RemoteId::Path(p) => p.name().to_string(),
_ => String::new(),
};
self.bodies
.get(&name)
.cloned()
.ok_or(RemoteError::NotFound(name))
}
}
}
async fn put(
&self,
_p: &RemotePath,
_b: Vec<u8>,
_pc: Option<dr_sync::Precondition>,
) -> Result<dr_sync::Validator, RemoteError> {
self.puts.fetch_add(1, Ordering::SeqCst);
*self.last_put.lock().unwrap() = _b;
Ok(dr_sync::Validator::new("v"))
}
async fn delete(
&self,
_id: &RemoteId,
_pc: Option<dr_sync::Precondition>,
) -> Result<(), RemoteError> {
Ok(())
}
async fn move_to(&self, f: &RemoteId, t: &RemotePath) -> Result<(), RemoteError> {
let from = match f {
RemoteId::Path(p) => p.name().to_string(),
_ => String::new(),
};
self.moves
.lock()
.unwrap()
.push((from, t.name().to_string()));
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<RemoteError>, 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);
}
// --- a damaged copy on the server (FR-NC-9) --------------------------
//
// The one exception to everything above. A file SQLite calls malformed is
// not a read that failed; it is a file no device will ever read again, and
// leaving it alone pins it there for every client at once.
/// The bytes of a catalog holding one collection, as a peer would upload.
fn a_catalog_with(collection: &str, name: &str) -> Vec<u8> {
let dir = std::env::temp_dir().join(format!("dr-generation-{name}"));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("catalog.sqlite");
let cat = dr_catalog::Catalog::open(&path).unwrap();
cat.connection()
.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES (?1, ?2, 0, 0, 1, 1)",
rusqlite::params![format!("uuid-{collection}"), collection],
)
.unwrap();
let snap = dir.join("snap.sqlite");
cat.snapshot_for_upload(&snap).unwrap();
let bytes = std::fs::read(&snap).unwrap();
let _ = std::fs::remove_dir_all(&dir);
bytes
}
/// A body that is definitely not a SQLite database as the current copy,
/// with a chosen advertised size and whatever generations the test wants.
async fn corrupt_remote(
body_len: usize,
advertised: Option<u64>,
generations: &[(usize, Vec<u8>)],
name: &str,
) -> (usize, SyncReport, Vec<(String, String)>) {
let (catalog_path, scratch) = fixture(name);
let (mut backend, puts, _) = Fussy::serving(None, Vec::new());
backend.advertise = advertised;
backend
.bodies
.insert(CATALOG_NAME.to_string(), vec![0xAB; body_len]);
for (n, bytes) in generations {
backend.bodies.insert(generation_name(*n), bytes.clone());
}
let moves = backend.moves.clone();
let mut report = SyncReport::default();
sync_catalog(
&backend,
&RemotePath::new(".darkroom-derived"),
&catalog_path,
&scratch,
&mut report,
)
.await
.unwrap();
let moves = moves.lock().unwrap().clone();
(puts.load(Ordering::SeqCst), report, moves)
}
#[tokio::test]
async fn a_damaged_catalog_that_arrived_whole_is_replaced() {
// The deadlock this exists to break: every device downloads the same
// unreadable file, every device declines to overwrite it, and
// collections and people stop crossing between devices for ever.
let (puts, report, moves) = corrupt_remote(64, Some(64), &[], "corrupt-whole").await;
assert_eq!(puts, 1, "ours goes up, to the staging name");
assert!(report.catalog_replaced, "and the report says what happened");
assert!(report.catalog_uploaded);
assert!(
!report.catalog_merged,
"there was nothing readable to merge"
);
// The damaged file is kept by the rotation, not thrown away.
assert!(moves.contains(&(CATALOG_NAME.into(), generation_name(1))));
assert_eq!(
moves.last().unwrap(),
&(UPLOAD_NAME.to_string(), CATALOG_NAME.to_string()),
"and the upload is moved into place last"
);
}
#[tokio::test]
async fn a_damaged_catalog_falls_back_to_the_generation_before_it() {
// What generations are for. The device that pushed the damaged copy
// may have been the only one holding some collection; the generation
// before it still has everything every device had agreed on.
let older = a_catalog_with("Iceland", "gen1");
let (puts, report, _) =
corrupt_remote(64, Some(64), &[(1, older)], "corrupt-with-gen").await;
assert!(report.catalog_merged, "merged from catalog.1.sqlite");
assert_eq!(report.collections_gained, 1, "and gained what it held");
assert!(report.catalog_replaced);
assert_eq!(puts, 1);
}
#[tokio::test]
async fn a_damaged_generation_is_skipped_for_the_one_behind_it() {
// Two bad pushes in a row must not be worse than one.
let older = a_catalog_with("Faroe", "gen2");
let (_, report, _) = corrupt_remote(
64,
Some(64),
&[(1, vec![0xCD; 64]), (2, older)],
"corrupt-two-deep",
)
.await;
assert!(report.catalog_merged);
assert_eq!(report.collections_gained, 1);
}
#[tokio::test]
async fn a_short_download_is_a_truncated_transfer_not_a_damaged_file() {
// The failure that matters most to get right. A phone on a flaky link
// aborts bodies constantly, and a partial download will not open
// either — treating that as damage would let one bad connection
// destroy a catalog the server was holding perfectly well.
let (puts, report, moves) = corrupt_remote(64, Some(4096), &[], "corrupt-short").await;
assert_eq!(puts, 0, "nothing is written over a copy that arrived short");
assert!(moves.is_empty(), "and nothing is rotated");
assert!(!report.catalog_replaced);
assert!(!report.catalog_uploaded);
}
#[tokio::test]
async fn a_size_the_server_will_not_confirm_leaves_the_copy_alone() {
// Not knowing is not the same as knowing it is whole. Without a size
// to compare against there is no way to tell damage from truncation,
// and the answer to "I cannot tell" has to stay "do not overwrite".
let (puts, report, _) = corrupt_remote(64, None, &[], "corrupt-unconfirmed").await;
assert_eq!(puts, 0);
assert!(!report.catalog_replaced);
}
#[tokio::test]
async fn a_push_rotates_oldest_first_and_lands_last() {
// The order is the safety: every destination is empty when it is
// moved into, so a failure at any step leaves a gap and never a
// missing current copy.
let (puts, report, moves) =
run_with_moves(Some(RemoteError::NotFound("nope".into())), "rotation").await;
assert_eq!(puts, 1);
assert!(report.catalog_uploaded);
assert_eq!(
moves,
vec![
(generation_name(2), generation_name(3)),
(generation_name(1), generation_name(2)),
(CATALOG_NAME.to_string(), generation_name(1)),
(UPLOAD_NAME.to_string(), CATALOG_NAME.to_string()),
]
);
}
#[tokio::test]
async fn an_upload_the_server_holds_at_the_wrong_size_is_not_rotated_in() {
// The assembly went wrong on the server. Rotating that into place would
// hand every other device the damaged file the generations exist to
// survive; discarding it costs one retry.
let (catalog_path, scratch) = fixture("wrong-size");
let (mut backend, puts) = Fussy::reading(Some(RemoteError::NotFound("nope".into())));
backend.staged_size = Some(7);
let moves = backend.moves.clone();
let mut report = SyncReport::default();
sync_catalog(
&backend,
&RemotePath::new(".darkroom-derived"),
&catalog_path,
&scratch,
&mut report,
)
.await
.unwrap();
assert_eq!(puts.load(Ordering::SeqCst), 1, "it was uploaded");
assert!(!report.catalog_uploaded, "but not accepted");
assert!(moves.lock().unwrap().is_empty(), "and nothing was rotated");
}
async fn run_with_moves(
fail_with: Option<RemoteError>,
name: &str,
) -> (usize, SyncReport, Vec<(String, String)>) {
let (catalog_path, scratch) = fixture(name);
let (backend, puts) = Fussy::reading(fail_with);
let moves = backend.moves.clone();
let mut report = SyncReport::default();
sync_catalog(
&backend,
&RemotePath::new(".darkroom-derived"),
&catalog_path,
&scratch,
&mut report,
)
.await
.unwrap();
let moves = moves.lock().unwrap().clone();
(puts.load(Ordering::SeqCst), report, moves)
}
// --- the place (FR-UI-8) ---------------------------------------------
//
// The same "do not write over what you could not read" rule as above, for a
// file with different stakes. A place is replaced rather than merged, so an
// unread copy costs no data — but it may be the *newer* one, and pushing
// ours over it would move the other device's photographer without ever
// having seen where they were.
fn a_place(at: i64, device: &str) -> dr_types::Place {
dr_types::Place {
at,
device: device.into(),
path: "2019/a.CR2".into(),
..Default::default()
}
}
/// A local place file with a chosen record already in it, or none.
fn local(name: &str, ours: Option<&dr_types::Place>) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!("dr-place-guard-{name}"));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("place.json");
if let Some(p) = ours {
crate::place::PlaceStore::open_at(path.clone()).save(p);
}
path
}
async fn exchange(
fail_with: Option<RemoteError>,
theirs: Option<&dr_types::Place>,
ours: Option<&dr_types::Place>,
name: &str,
) -> (usize, SyncReport, Option<dr_types::Place>, Vec<u8>) {
let body = theirs
.map(|p| crate::place::to_bytes(p).unwrap())
.unwrap_or_default();
let (backend, puts, last) = Fussy::serving(fail_with, body);
let path = local(name, ours);
let mut report = SyncReport::default();
sync_place(
&backend,
&RemotePath::new(".darkroom-derived"),
&path,
&mut report,
)
.await;
let after = crate::place::PlaceStore::open_at(path.clone()).load();
let sent = last.lock().unwrap().clone();
let _ = std::fs::remove_dir_all(path.parent().unwrap());
(puts.load(Ordering::SeqCst), report, after, sent)
}
/// A small valid profile, as bytes, named for `model`.
fn profile_bytes(model: &str) -> Vec<u8> {
dr_decode::dcp::Dcp {
name: "Test Standard".into(),
unique_camera_model: Some(model.into()),
copyright: None,
calibration_signature: None,
embed_policy: 0,
illuminants: [Some(21), None],
color_matrix: [
Some([[1.0, 0.0, 0.0], [0.0, 1.0, 0.0], [0.0, 0.0, 1.0]]),
None,
],
forward_matrix: [None, None],
hue_sat: [None, None],
look: dr_types::HueSatTable::new(2, 2, 1, false, vec![[5.0, 1.1, 1.0]; 4]),
tone_curve: None,
baseline_exposure_offset: 0.0,
}
.to_bytes()
.unwrap()
}
#[tokio::test]
async fn camera_profiles_travel_both_ways_and_once() {
// TRACES: FR-DEV-3e
// camera-profiles.md §13: ours goes up, theirs comes down, and a
// second pass with nothing new moves nothing.
let root = std::env::temp_dir().join(format!("dr-profile-sync-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let server = root.join("server");
let device = root.join("device");
std::fs::create_dir_all(server.join(".darkroom-derived/profiles")).unwrap();
std::fs::create_dir_all(&device).unwrap();
std::fs::write(device.join("Ours A.dcp"), profile_bytes("Ours A")).unwrap();
std::fs::write(
server.join(".darkroom-derived/profiles/Theirs B.dcp"),
profile_bytes("Theirs B"),
)
.unwrap();
// Not a profile, and must be left alone in both directions.
std::fs::write(server.join(".darkroom-derived/profiles/notes.txt"), b"x").unwrap();
std::fs::write(
server.join(".darkroom-derived/profiles/Broken C.dcp"),
b"not a profile",
)
.unwrap();
let backend = dr_sync_folder::FolderBackend::new(&server).unwrap();
let base = RemotePath::new(".darkroom-derived");
let mut report = SyncReport::default();
sync_profiles(&backend, &base, &device, &mut report).await;
assert_eq!(report.profiles_uploaded, 1);
assert_eq!(report.profiles_downloaded, 1, "the broken one is refused");
assert!(server
.join(".darkroom-derived/profiles/Ours A.dcp")
.exists());
assert!(device.join("Theirs B.dcp").exists());
assert!(!device.join("Broken C.dcp").exists());
assert!(!device.join("notes.txt").exists());
assert!(
dr_decode::dcp::find(Some("Theirs B"), "", "").is_some(),
"a fetched profile is matched on the next decode"
);
let mut again = SyncReport::default();
sync_profiles(&backend, &base, &device, &mut again).await;
assert_eq!((again.profiles_uploaded, again.profiles_downloaded), (0, 0));
let _ = std::fs::remove_dir_all(&root);
}
/// One device's preset files, under `root`.
fn preset_device(root: &Path, name: &str) -> PresetFiles {
PresetFiles {
library: root.join(name).join("presets.drpl"),
base: root.join(name).join("presets.base.drpl"),
}
}
fn preset_names(files: &PresetFiles) -> Vec<String> {
crate::preset_store::PresetStore::open_at(files.library.clone())
.load()
.names()
.map(str::to_string)
.collect()
}
fn save_presets(files: &PresetFiles, names: &[&str]) {
let mut library = dr_pipeline::PresetLibrary::default();
for name in names {
library
.insert(name, dr_pipeline::Preset::default())
.unwrap();
}
crate::preset_store::PresetStore::open_at(files.library.clone())
.save(&library)
.unwrap();
}
#[tokio::test]
async fn presets_reach_every_device_and_so_do_their_deletions() {
// TRACES: FR-DEV-6
let root = std::env::temp_dir().join(format!("dr-preset-sync-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let server = root.join("server");
std::fs::create_dir_all(server.join(".darkroom-derived")).unwrap();
let backend = dr_sync_folder::FolderBackend::new(&server).unwrap();
let base = RemotePath::new(".darkroom-derived");
let (desk, tablet) = (preset_device(&root, "desk"), preset_device(&root, "tablet"));
save_presets(&desk, &["Warm"]);
save_presets(&tablet, &["Mono"]);
for files in [&desk, &tablet, &desk] {
sync_presets(&backend, &base, files, &mut SyncReport::default()).await;
}
assert_eq!(preset_names(&desk), vec!["Mono", "Warm"]);
assert_eq!(preset_names(&tablet), vec!["Mono", "Warm"]);
// Deleted on the desk: gone from the tablet, not back on the desk.
save_presets(&desk, &["Mono"]);
for files in [&desk, &tablet, &desk] {
sync_presets(&backend, &base, files, &mut SyncReport::default()).await;
}
assert_eq!(preset_names(&desk), vec!["Mono"]);
assert_eq!(preset_names(&tablet), vec!["Mono"]);
// Settled: a pass with nothing new writes nothing.
let mut quiet = SyncReport::default();
sync_presets(&backend, &base, &tablet, &mut quiet).await;
assert!(!quiet.presets_uploaded && !quiet.presets_adopted);
let _ = std::fs::remove_dir_all(&root);
}
#[tokio::test]
async fn a_preset_library_the_server_cannot_read_is_left_alone() {
// TRACES: FR-DEV-6
// A newer build's format reads as unreadable here, and is that build's
// presets: neither written over nor taken as an empty library.
let root = std::env::temp_dir().join(format!("dr-preset-unread-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let server = root.join("server");
let held = server.join(".darkroom-derived/presets/library.drpl");
std::fs::create_dir_all(held.parent().unwrap()).unwrap();
std::fs::write(&held, "drpl 9999\n").unwrap();
let backend = dr_sync_folder::FolderBackend::new(&server).unwrap();
let desk = preset_device(&root, "desk");
save_presets(&desk, &["Warm"]);
let mut report = SyncReport::default();
sync_presets(
&backend,
&RemotePath::new(".darkroom-derived"),
&desk,
&mut report,
)
.await;
assert!(!report.presets_uploaded);
assert_eq!(std::fs::read_to_string(&held).unwrap(), "drpl 9999\n");
assert_eq!(preset_names(&desk), vec!["Warm"]);
let _ = std::fs::remove_dir_all(&root);
}
#[tokio::test]
async fn a_place_that_could_not_be_read_is_never_written_over() {
// A dehydrated placeholder, and the record on the server may well be
// newer than ours. Uploading blind would discard it.
let ours = a_place(100, "desktop");
let (puts, report, after, _) = exchange(
Some(RemoteError::NotMaterialised("place.json".into())),
None,
Some(&ours),
"notmaterialised",
)
.await;
assert_eq!(puts, 0, "must not upload over a place it could not read");
assert!(!report.place_uploaded);
assert!(!report.place_adopted);
assert_eq!(after.unwrap().at, 100, "and ours is untouched");
}
#[tokio::test]
async fn a_first_exchange_uploads_ours() {
// Genuinely nothing there: ours is the whole truth, and refusing to
// push it would mean the place never travels at all.
let ours = a_place(100, "desktop");
let (puts, report, _, sent) = exchange(
Some(RemoteError::NotFound("nope".into())),
None,
Some(&ours),
"firstrun",
)
.await;
assert_eq!(puts, 1);
assert!(report.place_uploaded);
assert_eq!(
crate::place::from_bytes(&sent, "sent").unwrap().device,
"desktop"
);
}
#[tokio::test]
async fn a_newer_record_from_another_device_is_adopted_and_not_pushed_back() {
// The handover, and the half that stops it oscillating: having taken
// theirs, ours *is* theirs, so there is nothing left to send.
let ours = a_place(100, "desktop");
let theirs = a_place(200, "tablet");
let (puts, report, after, _) = exchange(None, Some(&theirs), Some(&ours), "adopt").await;
assert!(report.place_adopted);
assert_eq!(after.unwrap().device, "tablet");
assert_eq!(puts, 0, "nothing to say that the server does not know");
assert!(!report.place_uploaded);
}
#[tokio::test]
async fn an_older_record_on_the_server_is_replaced_by_ours() {
let ours = a_place(300, "desktop");
let theirs = a_place(200, "tablet");
let (puts, report, after, sent) = exchange(None, Some(&theirs), Some(&ours), "push").await;
assert!(!report.place_adopted);
assert_eq!(after.unwrap().device, "desktop", "ours stands");
assert_eq!(puts, 1);
assert!(report.place_uploaded);
assert_eq!(crate::place::from_bytes(&sent, "sent").unwrap().at, 300);
}
#[tokio::test]
async fn two_idle_devices_settle_rather_than_ping_pong() {
// Equal timestamps mean equal records. A pass that decided it had
// something to say here would put a PUT on every sync of every device,
// for ever.
let same = a_place(100, "desktop");
let (puts, report, _, _) = exchange(None, Some(&same), Some(&same), "settle").await;
assert_eq!(puts, 0);
assert!(!report.place_adopted);
assert!(!report.place_uploaded);
}
#[tokio::test]
async fn a_device_with_no_place_of_its_own_takes_theirs() {
// A second machine signing into an existing library. There is nothing
// local to compare against, so anything on the server wins.
let theirs = a_place(200, "tablet");
let (puts, report, after, _) = exchange(None, Some(&theirs), None, "fresh").await;
assert!(report.place_adopted);
assert_eq!(after.unwrap().device, "tablet");
assert_eq!(puts, 0);
}
#[tokio::test]
async fn a_torn_record_on_the_server_is_replaced_rather_than_obeyed() {
// Unlike the catalog, an unparseable place is not a reason to hold
// back: there is nothing inside it to lose, and leaving a corrupt file
// in place would mean the exchange never recovers.
let ours = a_place(100, "desktop");
let (backend, puts, _) = Fussy::serving(None, b"{\"at\": 12, \"scr".to_vec());
let path = local("torn", Some(&ours));
let mut report = SyncReport::default();
sync_place(
&backend,
&RemotePath::new(".darkroom-derived"),
&path,
&mut report,
)
.await;
assert_eq!(puts.load(Ordering::SeqCst), 1, "ours goes over the rubble");
assert!(report.place_uploaded);
let _ = std::fs::remove_dir_all(path.parent().unwrap());
}
}