Files
DarkRoom/ui/dr-ui/src/derived_sync.rs
T
dtourolle f100db89ca Verify the catalog snapshot before it is sent, and after it lands
Two checks around the upload, both cheap next to what they prevent.

Before: the snapshot is quick_checked before it leaves. It is the copy
every other device merges from, and a damaged one costs each of them a
download, a failed merge and a refusal to push.

After: the staged upload's size on the server is compared to the bytes
sent before it is rotated into place. A chunked upload is assembled
server-side, and an assembly that goes wrong is a file of plausible
size no device can open — caught here, on the device that caused it,
for one listing; otherwise on every other device, after the fact. A
mismatch, or a size the server will not confirm, discards the upload
and leaves the current copy and its generations untouched.
2026-09-13 19:31:58 +02:00

1821 lines
73 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 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,
// 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,
}
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
}
}
/// Progress from the sync worker.
#[derive(Debug)]
pub enum SyncMessage {
Status(String),
Finished(Box<SyncReport>),
Failed(String),
}
/// Push shards and the catalog, and take anything newer from the server.
///
/// Runs on its own thread with its own runtime, like every other network path
/// here — the Slint loop must never block (NFR-P9).
pub fn spawn_sync(
conn: Connection,
root: String,
thumbs_dir: PathBuf,
catalog_path: PathBuf,
place_path: PathBuf,
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();
std::thread::spawn(move || {
// A current-thread runtime here is what made the library scan itself
// rather than adopt the shards the server already had; see net_runtime.
let rt = match crate::net_runtime::build() {
Ok(rt) => rt,
Err(e) => {
let _ = tx.send(SyncMessage::Failed(e.to_string()));
return;
}
};
rt.block_on(async {
let backend = match crate::remote::connect(&conn) {
Ok(b) => b,
Err(e) => {
let _ = tx.send(SyncMessage::Failed(e.to_string()));
return;
}
};
match run(
&*backend,
&root,
&thumbs_dir,
&catalog_path,
&place_path,
&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
}
// Eight, because a sync touches eight distinct things — the same reason
// `spawn_face_sweep` 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,
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;
let _ = tx.send(SyncMessage::Status("checking thumbnails…".into()));
sync_shards(backend, &base, thumbs_dir, scratch, &mut report).await?;
let _ = tx.send(SyncMessage::Status("checking faces…".into()));
sync_face_shards(
backend,
&base,
catalog_path,
scratch,
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?;
// 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 {
let path = store.shard_path(shard.id);
let Ok(bytes) = std::fs::read(&path) else {
continue;
};
let name = shard_name(&client, shard.id);
// **Size, sealed or not.** This used to take the server merely *having*
// a sealed shard as proof it had the whole thing — "byte-identical by
// construction". It is not: an open shard is uploaded on every pass as
// it grows, so the server routinely holds a partial copy of a shard
// that is later filled and sealed. From the moment of sealing, the name
// matched and the size was never looked at again, and the completed
// shard could never be sent. The client id in the name still means
// nobody else can have written it, so size is a sound comparison.
let skip = match remote.get(&name) {
Some(size) => *size == bytes.len() as u64,
None => false,
};
if skip {
continue;
}
let target = RemotePath::new(format!("{}/{name}", base.as_str()));
match backend.put(&target, bytes, None).await {
Ok(_) => report.shards_uploaded += 1,
// One shard failing must not abort the rest: they are independent
// and the next pass retries.
Err(e) => log::warn!("uploading {name}: {e}"),
}
}
// ---- download --------------------------------------------------------
let mut store = store;
for (name, size) in &remote {
let Some((owner, id)) = parse_shard(name) else {
continue;
};
if owner == client {
continue;
}
// Merged already, at the size it still has. A sealed shard never
// reaches here twice; a peer's open one does each time it grows, which
// is what carries its later thumbnails across.
if store.adopted(name) == Some(*size) {
continue;
}
if owner.is_empty() && legacy_upload_of_ours(&store, id, *size) {
let _ = store.record_adopted(name, *size);
continue;
}
let source = RemotePath::new(format!("{}/{name}", base.as_str()));
// Fetched where it is only a placeholder: a shard that will not open
// is a peer's thumbnails never merging, and on a library the client
// keeps dehydrated that would be every shard, every pass, silently.
let bytes = match read_derived(backend, &source).await {
Ok(b) => b,
Err(e) => {
log::warn!("downloading {name}: {e}");
continue;
}
};
// Written to scratch and merged, rather than dropped into the store
// directory: a downloaded shard's *id* is the other device's numbering,
// and two devices independently fill shard 0.
let tmp = scratch.join(name);
if std::fs::write(&tmp, &bytes).is_err() {
continue;
}
match store.merge_shard(&tmp) {
Ok(n) => {
report.shards_downloaded += 1;
report.thumbnails_adopted += n;
// Recorded only on success, so a failed merge is retried next
// pass rather than written off.
let _ = store.record_adopted(name, *size);
}
Err(e) => log::warn!("merging {name}: {e}"),
}
let _ = std::fs::remove_file(&tmp);
}
Ok(())
}
/// Push and pull face shards, so a second device does not re-index the library.
///
/// Mirrors [`sync_shards`] deliberately, down to the sealed-shard skip and the
/// adopted ledger: face data has exactly the properties that made that design
/// right for thumbnails. It is bulk, it is immutable once written, and it is
/// byte-identical on every device, because the same model over the same proxy
/// is deterministic.
///
/// The catalog is the source and the destination; the shards are only the
/// carrier. So this exports the catalog's new faces into the local shard store
/// first, syncs the shards, and imports whatever arrived back into the catalog.
async fn sync_face_shards(
backend: &dyn RemoteBackend,
base: &RemotePath,
catalog_path: &Path,
scratch: &Path,
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;
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-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();
std::thread::spawn(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)
}
#[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());
}
}