The catalog sync refuses to upload when it cannot read the server's copy, because the upload is a read-modify-write and writing blind would discard another device's collections. That is the right rule for a timeout, a dropped connection or a newer schema — the remote is fine, only our view of it failed. A file SQLite calls malformed is not that. No device will ever read it again, so refusing to write over it preserves nothing — and every client declines in turn, pinning the damaged file in place for good. Collections and people then stop crossing between devices on all of them at once, each logging "catalog not pushed" on every pass. This library did exactly that from 2026-09-07, on the desktop and on a freshly installed phone alike, while 32 collections sat undelivered. Now a copy that arrived whole and still will not open is set aside under a dated name and replaced by ours. Whole is checked against the size the server advertises: a truncated download will not open either, and on a phone that is the far likelier story, so anything short — or any size the listing cannot confirm — is treated as the transport failure it is and the server's copy is left alone. A placeholder's size is not trusted for the comparison, since it means nothing. The report says when this happened, and the log line calls it "pushed over a damaged copy" rather than folding it into an ordinary push: it is the one push that discarded something.
1516 lines
60 KiB
Rust
1516 lines
60 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.
|
|
///
|
|
/// Only collections merge — see [`dr_catalog::sync`]. The rest of a catalog
|
|
/// describes local state (folder ETags, cache paths, job rows) and importing
|
|
/// another device's version would be actively wrong.
|
|
async fn sync_catalog(
|
|
backend: &dyn RemoteBackend,
|
|
base: &RemotePath,
|
|
catalog_path: &Path,
|
|
scratch: &Path,
|
|
report: &mut SyncReport,
|
|
) -> Result<(), String> {
|
|
let remote_name = "catalog.sqlite";
|
|
let target = RemotePath::new(format!("{}/{remote_name}", base.as_str()));
|
|
|
|
// ---- take theirs first -----------------------------------------------
|
|
//
|
|
// Merging before uploading means our upload carries the union rather than
|
|
// only our own half, so a third device syncing next gets everything in one
|
|
// fetch.
|
|
// TRACES: FR-NC-9 | FR-NC-6c
|
|
// A read that fails for any reason other than "there is not one yet" must
|
|
// stop the upload below. This is a read-modify-write over a file another
|
|
// device also writes, so skipping the read does not merely lose an
|
|
// optimisation — it turns the write into a clobber, and the other device's
|
|
// collections and their members go with it.
|
|
//
|
|
// The shape was previously `if let Ok(bytes) = ...`, which swallowed every
|
|
// failure into "no remote catalog" and carried straight on to the upload.
|
|
let theirs = match read_derived(backend, &target).await {
|
|
Ok(bytes) => Some(bytes),
|
|
// Genuinely the first sync of this library. Nothing to merge, and
|
|
// ours is the whole truth.
|
|
Err(RemoteError::NotFound(_)) => None,
|
|
Err(e) => {
|
|
log::warn!(
|
|
"not pushing the catalog: the copy on the server could not be read ({e}); \
|
|
uploading over it would discard whatever another device put there"
|
|
);
|
|
return Ok(());
|
|
}
|
|
};
|
|
|
|
if let Some(bytes) = theirs {
|
|
let downloaded = scratch.join("catalog-remote.sqlite");
|
|
if std::fs::write(&downloaded, &bytes).is_ok() {
|
|
match dr_catalog::Catalog::open(catalog_path) {
|
|
Ok(catalog) => match catalog.merge_remote_catalog(&downloaded) {
|
|
Ok(merge) => {
|
|
report.catalog_merged = true;
|
|
report.collections_gained = merge.inserted + merge.updated;
|
|
report.members_gained = merge.members_added;
|
|
}
|
|
// TRACES: FR-NC-9
|
|
// Damaged beyond reading, and the whole file arrived. See
|
|
// [`replace_corrupt_remote`] for why this one failure is
|
|
// the exception to "never write over what you could not
|
|
// read": there is nothing left in it to preserve, and
|
|
// refusing for ever is what pinned a damaged file in place
|
|
// on every device for a week.
|
|
Err(dr_catalog::CatalogError::Corrupt { detail }) => {
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
if !replace_corrupt_remote(backend, base, remote_name, &bytes, &detail)
|
|
.await
|
|
{
|
|
return Ok(());
|
|
}
|
|
report.catalog_replaced = true;
|
|
}
|
|
// Unreadable is not the same as absent: it may be a newer
|
|
// format, or a torn upload. Ours must not go over it.
|
|
Err(e) => {
|
|
log::warn!("not pushing the catalog: merging the server's copy: {e}");
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
return Ok(());
|
|
}
|
|
},
|
|
Err(e) => {
|
|
log::warn!("not pushing the catalog: opening ours to merge: {e}");
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
return Ok(());
|
|
}
|
|
}
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
}
|
|
}
|
|
|
|
// ---- then push ours --------------------------------------------------
|
|
//
|
|
// Never the live file: committed transactions can sit in the `-wal` with
|
|
// the main file lagging, so copying it uploads a torn snapshot. The backup
|
|
// API serialises against writers instead of racing them.
|
|
let snapshot = scratch.join("catalog-upload.sqlite");
|
|
let catalog = dr_catalog::Catalog::open(catalog_path).map_err(|e| e.to_string())?;
|
|
catalog
|
|
.snapshot_for_upload(&snapshot)
|
|
.map_err(|e| e.to_string())?;
|
|
|
|
let bytes = std::fs::read(&snapshot).map_err(|e| e.to_string())?;
|
|
match backend.put(&target, bytes, None).await {
|
|
Ok(_) => report.catalog_uploaded = true,
|
|
Err(e) => log::warn!("uploading catalog: {e}"),
|
|
}
|
|
let _ = std::fs::remove_file(&snapshot);
|
|
|
|
Ok(())
|
|
}
|
|
|
|
/// TRACES: FR-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
|
|
/// Set a damaged remote catalog aside so ours can replace it.
|
|
///
|
|
/// # Why this is allowed to destroy something
|
|
///
|
|
/// Everything else in [`sync_catalog`] refuses to upload when it could not read
|
|
/// the server's copy, because the upload is a read-modify-write and writing
|
|
/// blind discards another device's collections. That rule is right for every
|
|
/// failure it was written for — a timeout, a dropped connection, a newer schema
|
|
/// — because in all of them the remote is *fine* and only our view of it
|
|
/// failed.
|
|
///
|
|
/// A file SQLite calls malformed is not that. It is not a view that failed; it
|
|
/// is a file whose contents no device can ever read again. Refusing to write
|
|
/// over it preserves nothing, and every client then declines in turn: the
|
|
/// damaged copy is pinned in place for ever, and collections and people stop
|
|
/// crossing between devices silently, on all of them at once. That is not
|
|
/// theoretical — it is what this library did from 2026-09-07, on the desktop
|
|
/// and on a freshly installed phone alike, both logging "catalog not pushed"
|
|
/// on every pass for a week while 32 collections sat undelivered.
|
|
///
|
|
/// # What makes it safe
|
|
///
|
|
/// **The whole file has to have arrived.** A truncated download is also
|
|
/// unreadable, and it is the far more likely story on a phone — this library's
|
|
/// logs are full of aborted bodies and DNS failures. So the size the server
|
|
/// advertises is compared against what we actually received, and anything short
|
|
/// is treated as the transport failure it is. Only a complete file that still
|
|
/// will not open is judged damaged.
|
|
///
|
|
/// **Nothing is deleted.** The damaged bytes are uploaded beside the catalog
|
|
/// under a dated name first, and the replacement only proceeds once that has
|
|
/// landed. If some later build learns to salvage collections out of a damaged
|
|
/// catalog, the file is still there to salvage them from.
|
|
///
|
|
/// Returns whether the caller may now push over the remote.
|
|
async fn replace_corrupt_remote(
|
|
backend: &dyn RemoteBackend,
|
|
base: &RemotePath,
|
|
remote_name: &str,
|
|
bytes: &[u8],
|
|
detail: &str,
|
|
) -> bool {
|
|
let advertised = match remote_size(backend, base, remote_name).await {
|
|
Some(n) => n,
|
|
None => {
|
|
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 != bytes.len() as u64 {
|
|
log::warn!(
|
|
"not pushing the catalog: the copy on the server will not open ({detail}), \
|
|
but only {} of {advertised} bytes arrived — that is a truncated download, \
|
|
not a damaged file, so the server's copy is left alone",
|
|
bytes.len()
|
|
);
|
|
return false;
|
|
}
|
|
|
|
let stamp = std::time::SystemTime::now()
|
|
.duration_since(std::time::UNIX_EPOCH)
|
|
.map(|d| d.as_secs())
|
|
.unwrap_or(0);
|
|
let aside_name = format!("catalog.corrupt-{stamp}.sqlite");
|
|
let aside = RemotePath::new(format!("{}/{aside_name}", base.as_str()));
|
|
|
|
if let Err(e) = backend.put(&aside, bytes.to_vec(), None).await {
|
|
log::warn!(
|
|
"not pushing the catalog: the copy on the server is damaged ({detail}), \
|
|
but it could not be set aside as {aside_name} ({e}), and it will not be \
|
|
overwritten until a copy of it is safe"
|
|
);
|
|
return false;
|
|
}
|
|
|
|
log::warn!(
|
|
"the catalog on the server is damaged ({detail}) and all {advertised} bytes of it \
|
|
arrived, so no device can read it. Kept as {aside_name}; replacing it with this \
|
|
device's copy, which is what lets collections and people sync again"
|
|
);
|
|
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>,
|
|
}
|
|
|
|
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,
|
|
},
|
|
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 Some(size) = self.advertise else {
|
|
return Ok(Vec::new());
|
|
};
|
|
let path = RemotePath::new(format!("{}/catalog.sqlite", dir.as_str()));
|
|
Ok(vec![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,
|
|
}])
|
|
}
|
|
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 => Ok(self.body.clone()),
|
|
}
|
|
}
|
|
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> {
|
|
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.
|
|
|
|
/// A body that is definitely not a SQLite database, with a chosen size the
|
|
/// listing will or will not agree with.
|
|
async fn corrupt_remote(
|
|
body_len: usize,
|
|
advertised: Option<u64>,
|
|
name: &str,
|
|
) -> (usize, SyncReport, Vec<u8>) {
|
|
let (catalog_path, scratch) = fixture(name);
|
|
let (mut backend, puts, last) = Fussy::serving(None, vec![0xAB; body_len]);
|
|
backend.advertise = advertised;
|
|
let mut report = SyncReport::default();
|
|
sync_catalog(
|
|
&backend,
|
|
&RemotePath::new(".darkroom-derived"),
|
|
&catalog_path,
|
|
&scratch,
|
|
&mut report,
|
|
)
|
|
.await
|
|
.unwrap();
|
|
let sent = last.lock().unwrap().clone();
|
|
(puts.load(Ordering::SeqCst), report, sent)
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn a_damaged_catalog_that_arrived_whole_is_set_aside_and_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, _) = corrupt_remote(64, Some(64), "corrupt-whole").await;
|
|
assert_eq!(puts, 2, "the damaged copy is kept, then ours goes over it");
|
|
assert!(report.catalog_replaced, "and the report says what happened");
|
|
assert!(report.catalog_uploaded);
|
|
assert!(!report.catalog_merged, "there was nothing to merge");
|
|
}
|
|
|
|
#[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, _) = corrupt_remote(64, Some(4096), "corrupt-short").await;
|
|
assert_eq!(puts, 0, "nothing is written over a copy that arrived short");
|
|
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);
|
|
}
|
|
|
|
// --- 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());
|
|
}
|
|
}
|