Build and test / Desktop (Linux) (push) Successful in 2h5m41s
Build and test / Layer separation (push) Successful in 50s
🐳 Android image / Build and push (push) Successful in 2s
Build and test / android-image (push) Successful in 2s
Traceability / Requirement traces (push) Successful in 1m37s
Build and test / Android (aarch64) (push) Successful in 59m40s
The catalog clobber existed because `if let Ok(bytes)` had a failure arm that was never exercised. Fixing it without covering that arm leaves the next person free to collapse it back. Three tests against a backend whose read fails in a chosen way, counting writes — because what went wrong was not a wrong value but a write that should not have happened at all: - a dehydrated snapshot uploads nothing - an unreadable one uploads nothing either, since "refused" is no more "absent" than "not downloaded" is - and a genuine first sync still uploads, which is the half that keeps `NotFound` distinct from `NotMaterialised` rather than merely cautious Checked against the original shape: the first two fail on it and the third passes. A guard test that cannot tell the bug from the fix is decoration. `async-trait` joins dev-dependencies to stand the double up behind `dyn RemoteBackend`.
918 lines
36 KiB
Rust
918 lines
36 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,
|
|
pub collections_gained: usize,
|
|
/// Membership rows the merge brought in (FR-CAT-7).
|
|
///
|
|
/// Counted apart from `collections_gained`, which counts the collections
|
|
/// themselves. A sync that files 127 photographs into a collection both
|
|
/// devices already had gains *no* collection, and reporting only the
|
|
/// former left the sidebar showing no count beside a full collection.
|
|
pub members_gained: usize,
|
|
|
|
// Face data is counted apart from thumbnails for the same reason keywords
|
|
// are counted apart from collections: "adopted 4,812 faces" is a sentence
|
|
// the user can act on, and folding it into the thumbnail count would hide
|
|
// the one number that says whether this device still has hours of indexing
|
|
// ahead of it.
|
|
pub face_shards_uploaded: usize,
|
|
pub face_shards_downloaded: usize,
|
|
/// Images whose faces this device took from a peer instead of detecting.
|
|
pub faces_adopted: usize,
|
|
}
|
|
|
|
impl SyncReport {
|
|
pub fn did_anything(&self) -> bool {
|
|
self.shards_uploaded > 0
|
|
|| self.shards_downloaded > 0
|
|
|| self.catalog_uploaded
|
|
|| self.catalog_merged
|
|
|| self.face_shards_uploaded > 0
|
|
|| self.face_shards_downloaded > 0
|
|
}
|
|
}
|
|
|
|
/// Progress from the sync worker.
|
|
#[derive(Debug)]
|
|
pub enum SyncMessage {
|
|
Status(String),
|
|
Finished(Box<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,
|
|
scratch: PathBuf,
|
|
) -> 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, &scratch, &tx).await {
|
|
Ok(report) => {
|
|
let _ = tx.send(SyncMessage::Finished(Box::new(report)));
|
|
}
|
|
Err(e) => {
|
|
let _ = tx.send(SyncMessage::Failed(e));
|
|
}
|
|
}
|
|
});
|
|
});
|
|
|
|
rx
|
|
}
|
|
|
|
async fn run(
|
|
backend: &dyn RemoteBackend,
|
|
root: &str,
|
|
thumbs_dir: &Path,
|
|
catalog_path: &Path,
|
|
scratch: &Path,
|
|
tx: &std::sync::mpsc::Sender<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, &mut report, tx).await?;
|
|
|
|
let _ = tx.send(SyncMessage::Status("checking collections…".into()));
|
|
sync_catalog(backend, &base, catalog_path, scratch, &mut report).await?;
|
|
|
|
Ok(report)
|
|
}
|
|
|
|
/// The derived folder for a library root.
|
|
fn derived_path(root: &str) -> RemotePath {
|
|
if root.is_empty() {
|
|
RemotePath::new(DERIVED_DIR)
|
|
} else {
|
|
RemotePath::new(format!("{root}/{DERIVED_DIR}"))
|
|
}
|
|
}
|
|
|
|
/// Exchange thumbnail shards with the server.
|
|
///
|
|
/// Upload what the server lacks, download what we lack. Sealed shards are
|
|
/// immutable, so a name match is a content match and nothing needs comparing
|
|
/// beyond existence — which is what keeps a steady-state sync to one listing.
|
|
///
|
|
/// # Why the name carries a client
|
|
///
|
|
/// Shard ids are per-store: every client fills its own numbering from 0, so
|
|
/// "shard 3" names different thumbnails on each device. A flat `shard-0003`
|
|
/// remote namespace therefore has two clients writing one name — the second
|
|
/// upload overwrites content the first still believes is published — and
|
|
/// leaves a client no way to tell a peer's shard 3 from its own, so the only
|
|
/// safe reading of "I already have 3" is to skip it and never adopt anything.
|
|
/// Qualifying the name with [`ThumbStore::client_id`] gives each store its own
|
|
/// namespace, and [`ThumbStore::adopted`] then tracks what has been merged by
|
|
/// remote name instead of by our own ids.
|
|
async fn sync_shards(
|
|
backend: &dyn RemoteBackend,
|
|
base: &RemotePath,
|
|
thumbs_dir: &Path,
|
|
scratch: &Path,
|
|
report: &mut SyncReport,
|
|
) -> Result<(), String> {
|
|
let store = match ThumbStore::open(thumbs_dir) {
|
|
Ok(s) => s,
|
|
Err(e) => {
|
|
// A store that will not open can be neither read nor merged into,
|
|
// so there is no half of this worth attempting. It is not a sync
|
|
// failure: the next pass retries once the store is openable.
|
|
log::debug!("thumbnail store unavailable: {e}");
|
|
return Ok(());
|
|
}
|
|
};
|
|
let client = store.client_id().to_string();
|
|
|
|
let remote: std::collections::HashMap<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,
|
|
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 = crate::identity_ui::MODEL_ID;
|
|
|
|
// ---- everything this device has detected, into the shards ------------
|
|
if let Ok(catalog) = dr_catalog::Catalog::open(catalog_path) {
|
|
// The first export after a re-index walks the whole library and writes
|
|
// thousands of images. Saying so is the difference between a sync that
|
|
// looks slow and one that looks broken — this pass sat on "checking
|
|
// faces…" for eight minutes and was reported as a hang.
|
|
let mut say = |done: usize, total: usize| {
|
|
let _ = tx.send(SyncMessage::Status(format!(
|
|
"preparing faces to send: {done}/{total}"
|
|
)));
|
|
};
|
|
match face_shard::export_to_shards_reporting(
|
|
catalog.connection(),
|
|
&mut store,
|
|
model,
|
|
&mut say,
|
|
) {
|
|
Ok(0) => {}
|
|
Ok(n) => log::info!("face sync: {n} newly indexed image(s) ready to upload"),
|
|
Err(e) => log::warn!("face sync: exporting to shards: {e}"),
|
|
}
|
|
}
|
|
|
|
// Its own folder under the derived directory, so a client that does not
|
|
// care about faces lists thumbnails without paging past them.
|
|
let face_base = RemotePath::new(format!("{}/faces", base.as_str()));
|
|
let _ = backend.create_dir(&face_base).await;
|
|
|
|
let remote: std::collections::HashMap<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;
|
|
}
|
|
// Unreadable is not the same as absent: it may be a newer
|
|
// format, or a torn upload. Ours must not go over it.
|
|
Err(e) => {
|
|
log::warn!("not pushing the catalog: merging the server's copy: {e}");
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
return Ok(());
|
|
}
|
|
},
|
|
Err(e) => {
|
|
log::warn!("not pushing the catalog: opening ours to merge: {e}");
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
return Ok(());
|
|
}
|
|
}
|
|
let _ = std::fs::remove_file(&downloaded);
|
|
}
|
|
}
|
|
|
|
// ---- then push ours --------------------------------------------------
|
|
//
|
|
// Never the live file: committed transactions can sit in the `-wal` with
|
|
// the main file lagging, so copying it uploads a torn snapshot. The backup
|
|
// API serialises against writers instead of racing them.
|
|
let snapshot = scratch.join("catalog-upload.sqlite");
|
|
let catalog = dr_catalog::Catalog::open(catalog_path).map_err(|e| e.to_string())?;
|
|
catalog
|
|
.snapshot_for_upload(&snapshot)
|
|
.map_err(|e| e.to_string())?;
|
|
|
|
let bytes = std::fs::read(&snapshot).map_err(|e| e.to_string())?;
|
|
match backend.put(&target, bytes, None).await {
|
|
Ok(_) => report.catalog_uploaded = true,
|
|
Err(e) => log::warn!("uploading catalog: {e}"),
|
|
}
|
|
let _ = std::fs::remove_file(&snapshot);
|
|
|
|
Ok(())
|
|
}
|
|
|
|
/// TRACES: FR-NC-6c
|
|
/// Read a derived file, fetching its content first if only a placeholder is
|
|
/// here.
|
|
///
|
|
/// Derived state lives *inside the library folder*, so on a placeholder
|
|
/// library a sync client dehydrates a shard or a catalog snapshot exactly as
|
|
/// it dehydrates a photograph. Unlike a photograph, these are ours, and none of
|
|
/// them can be skipped: a shard that will not open is face data that never
|
|
/// merges, and a catalog snapshot that will not open is the other device's
|
|
/// collections.
|
|
///
|
|
/// So this fetches rather than giving up — and where it cannot, it says so
|
|
/// with the error rather than an empty result, because the callers below treat
|
|
/// "nothing there" as licence to write their own copy (ARCH §9.0a).
|
|
async fn read_derived(
|
|
backend: &dyn RemoteBackend,
|
|
path: &RemotePath,
|
|
) -> Result<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,
|
|
}
|
|
}
|
|
|
|
fn shard_name(client: &str, id: u32) -> String {
|
|
format!("shard-{client}-{id:04}.sqlite")
|
|
}
|
|
|
|
/// The client that wrote a remote shard and its id in that client's numbering,
|
|
/// or `None` if the name is not a shard.
|
|
///
|
|
/// Also accepts the flat `shard-NNNN.sqlite` written before names carried a
|
|
/// client, reporting an empty owner: those belong to nobody identifiable, so
|
|
/// they read as foreign and are adopted once like any peer's. Nothing is ever
|
|
/// uploaded under that form again.
|
|
///
|
|
/// Guards the download loop against adopting the catalog, a stray file, or
|
|
/// anything else the folder happens to contain.
|
|
fn parse_shard(name: &str) -> Option<(&str, u32)> {
|
|
let stem = name.strip_prefix("shard-")?.strip_suffix(".sqlite")?;
|
|
match stem.rsplit_once('-') {
|
|
Some((client, id)) => Some((client, id.parse().ok()?)),
|
|
None => Some(("", stem.parse().ok()?)),
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn derived_folder_sits_under_the_library_root() {
|
|
// Outside the root the account may not have write access — the root is
|
|
// the only thing the user granted.
|
|
assert_eq!(
|
|
derived_path("PhotosRaw").as_str(),
|
|
"PhotosRaw/.darkroom-derived"
|
|
);
|
|
// A library at the account root still gets a relative path.
|
|
assert_eq!(derived_path("").as_str(), ".darkroom-derived");
|
|
}
|
|
|
|
#[test]
|
|
fn shard_names_round_trip() {
|
|
assert_eq!(
|
|
shard_name("a1b2c3d4e5f6", 0),
|
|
"shard-a1b2c3d4e5f6-0000.sqlite"
|
|
);
|
|
assert_eq!(
|
|
shard_name("a1b2c3d4e5f6", 42),
|
|
"shard-a1b2c3d4e5f6-0042.sqlite"
|
|
);
|
|
assert_eq!(
|
|
parse_shard(&shard_name("a1b2c3d4e5f6", 7)),
|
|
Some(("a1b2c3d4e5f6", 7))
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn two_clients_shard_three_are_different_files() {
|
|
// The whole point: one client's numbering must not name another's
|
|
// shard, or the second upload overwrites the first's content and
|
|
// neither can tell the other's shards from its own.
|
|
assert_ne!(shard_name("aaaa", 3), shard_name("bbbb", 3));
|
|
assert_eq!(parse_shard(&shard_name("aaaa", 3)).unwrap().0, "aaaa");
|
|
assert_eq!(parse_shard(&shard_name("bbbb", 3)).unwrap().0, "bbbb");
|
|
}
|
|
|
|
#[test]
|
|
fn flat_names_read_as_belonging_to_nobody() {
|
|
// Written before the name carried a client. They must still parse, so
|
|
// a library synced by an older build is not stranded, and they must
|
|
// not match any live client id, so they are never mistaken for ours.
|
|
assert_eq!(parse_shard("shard-0042.sqlite"), Some(("", 42)));
|
|
assert_ne!(parse_shard("shard-0042.sqlite").unwrap().0, "a1b2c3d4e5f6");
|
|
}
|
|
|
|
#[test]
|
|
fn non_shard_files_are_not_adopted() {
|
|
// The folder also holds the catalog; downloading it as a shard would
|
|
// hand a catalog to the thumbnail merger.
|
|
assert_eq!(parse_shard("catalog.sqlite"), None);
|
|
assert_eq!(parse_shard("shard-0000.sqlite-wal"), None);
|
|
assert_eq!(parse_shard("notes.txt"), None);
|
|
assert_eq!(parse_shard("shard-abc.sqlite"), None);
|
|
assert_eq!(parse_shard("shard-a1b2c3-notanid.sqlite"), None);
|
|
}
|
|
|
|
#[test]
|
|
fn a_report_that_did_nothing_says_so() {
|
|
assert!(!SyncReport::default().did_anything());
|
|
assert!(SyncReport {
|
|
shards_uploaded: 1,
|
|
..Default::default()
|
|
}
|
|
.did_anything());
|
|
// Adopting thumbnails without moving a shard cannot happen, but the
|
|
// report must not claim work on collections alone either.
|
|
assert!(SyncReport {
|
|
catalog_merged: true,
|
|
..Default::default()
|
|
}
|
|
.did_anything());
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod catalog_guard_tests {
|
|
//! What `sync_catalog` does when it cannot read the server's copy.
|
|
//!
|
|
//! The bug these exist for was a control-flow one — `if let Ok(bytes)`
|
|
//! folding every failure into "there is none yet" and falling through to
|
|
//! the upload — so the thing to assert is not a value but *whether a write
|
|
//! happened at all*.
|
|
|
|
use super::*;
|
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
|
use std::sync::Arc;
|
|
|
|
/// A backend whose read fails in a chosen way, counting writes.
|
|
struct Fussy {
|
|
fail_with: Option<RemoteError>,
|
|
puts: Arc<AtomicUsize>,
|
|
caps: dr_sync::Capabilities,
|
|
}
|
|
|
|
impl Fussy {
|
|
fn reading(fail_with: Option<RemoteError>) -> (Self, Arc<AtomicUsize>) {
|
|
let puts = Arc::new(AtomicUsize::new(0));
|
|
(
|
|
Self {
|
|
fail_with,
|
|
puts: puts.clone(),
|
|
caps: dr_sync::Capabilities::minimal(),
|
|
},
|
|
puts,
|
|
)
|
|
}
|
|
}
|
|
|
|
#[async_trait::async_trait]
|
|
impl RemoteBackend for Fussy {
|
|
fn capabilities(&self) -> &dr_sync::Capabilities {
|
|
&self.caps
|
|
}
|
|
fn name(&self) -> &str {
|
|
"fussy"
|
|
}
|
|
async fn list(
|
|
&self,
|
|
_dir: &RemotePath,
|
|
_since: Option<&dr_sync::Validator>,
|
|
) -> Result<Vec<dr_sync::RemoteEntry>, RemoteError> {
|
|
Ok(Vec::new())
|
|
}
|
|
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(Vec::new()),
|
|
}
|
|
}
|
|
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);
|
|
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);
|
|
}
|
|
}
|