Merge: a catalog the server cannot damage for long, and backups on both sides

fix/corrupt-remote-catalog-deadlock. The catalog on the server had been
malformed since 7 September and every device declined to overwrite it,
so collections and people stopped syncing everywhere at once; the
damage came from two devices assembling chunks in one upload directory.
Each chunked upload now has its own directory, the server keeps three
generations of the catalog behind the current one, every push is
verified before and after, a damaged copy that arrived whole is merged
from the generation before it rather than pinned in place, and the
catalog is backed up daily as NFR-R2 always asked.
This commit is contained in:
2026-09-13 20:04:55 +02:00
6 changed files with 856 additions and 73 deletions
+91
View File
@@ -198,6 +198,54 @@ pub fn backup_before_migration(conn: &Connection, catalog: &Path) -> Result<(),
Ok(()) Ok(())
} }
/// How long a catalog may go without a backup before the next opportunity
/// takes one.
///
/// A day. The catalog is an index, so what a backup protects is the day's
/// worth of collection and people edits the sidecars do not hold — and a
/// second copy of a 130 MB file per launch would be a cost with nothing to
/// show for it when the user launches four times in an afternoon.
pub const BACKUP_EVERY: i64 = 24 * 60 * 60;
/// Whether [`BACKUP_EVERY`] has passed since the newest backup, or there is
/// none.
///
/// Read from the filenames, like [`backups`], so a restored or copied backup
/// directory answers the same way it did on the machine it came from.
pub fn backup_due(catalog: &Path) -> bool {
match backups(catalog).first() {
Some(newest) => now() - newest.taken_at >= BACKUP_EVERY,
None => true,
}
}
/// TRACES: NFR-R2
/// Take the scheduled backup, if one is due. Returns the file written, or
/// `None` when the newest is recent enough.
///
/// The scheduled half of NFR-R2 — the migration half is
/// [`backup_before_migration`]. "On a schedule" for an application that runs
/// when the user opens it means "at the next chance after a day has passed",
/// and the chance the caller picks is the end of a library sweep: the
/// catalog is quiet, the work is already off the UI thread, and it is the
/// moment a day's edits have just been consolidated.
///
/// A brand-new catalog with no images is not backed up: there is nothing in
/// it yet that a rescan would not rebuild, and the first backup would only be
/// a copy of an empty schema.
pub fn backup_if_due(conn: &Connection, catalog: &Path) -> Result<Option<PathBuf>, CatalogError> {
if !backup_due(catalog) {
return Ok(None);
}
let images: i64 = conn.query_row("SELECT count(*) FROM images", [], |r| r.get(0))?;
if images == 0 {
return Ok(None);
}
let path = backup(conn, catalog)?;
log::info!("scheduled backup of the catalog to {}", path.display());
Ok(Some(path))
}
/// The backups available for `catalog`, newest first. /// The backups available for `catalog`, newest first.
/// ///
/// Never fails: an unreadable or absent backup directory means there are no /// Never fails: an unreadable or absent backup directory means there are no
@@ -386,6 +434,49 @@ mod tests {
base base
} }
#[test]
fn a_scheduled_backup_is_taken_once_a_day_and_not_more() {
let dir = tempdir("scheduled");
let path = dir.join("catalog.sqlite");
fixture(&path, 3);
let cat = Catalog::open(&path).unwrap();
// Nothing yet: due.
assert!(backup_due(&path));
let first = backup_if_due(cat.connection(), &path).unwrap();
assert!(first.is_some(), "the first opportunity takes one");
// Taken just now: not due, and a second call does nothing.
assert!(!backup_due(&path));
assert_eq!(backup_if_due(cat.connection(), &path).unwrap(), None);
assert_eq!(backups(&path).len(), 1);
// Age the one backup past the interval by renaming it, since the
// timestamp is read from the name. Now it is due again.
let old = first.unwrap();
let aged = old
.parent()
.unwrap()
.join(format!("catalog-{}.sqlite", now() - BACKUP_EVERY - 1));
std::fs::rename(&old, &aged).unwrap();
assert!(backup_due(&path));
assert!(backup_if_due(cat.connection(), &path).unwrap().is_some());
assert_eq!(backups(&path).len(), 2);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn an_empty_catalog_is_not_worth_backing_up() {
let dir = tempdir("empty");
let path = dir.join("catalog.sqlite");
let cat = Catalog::open(&path).unwrap();
assert!(backup_due(&path), "due in principle");
assert_eq!(backup_if_due(cat.connection(), &path).unwrap(), None);
assert!(backups(&path).is_empty());
let _ = std::fs::remove_dir_all(&dir);
}
/// A catalog on disk with enough rows to span several pages, closed. /// A catalog on disk with enough rows to span several pages, closed.
/// ///
/// Closed matters: WAL means the rows are in `catalog.sqlite-wal` until /// Closed matters: WAL means the rows are in `catalog.sqlite-wal` until
+22
View File
@@ -54,9 +54,31 @@ pub fn checkpoint(conn: &Connection) -> Result<(), CatalogError> {
pub fn snapshot_for_upload(conn: &Connection, dest: &Path) -> Result<(), CatalogError> { pub fn snapshot_for_upload(conn: &Connection, dest: &Path) -> Result<(), CatalogError> {
let out = copy_to(conn, dest)?; let out = copy_to(conn, dest)?;
strip_face_crops(&out)?; strip_face_crops(&out)?;
verify_snapshot(&out)?;
Ok(()) Ok(())
} }
/// TRACES: NFR-R2
/// Refuse to hand over a snapshot that will not pass `quick_check`.
///
/// The upload is the copy every other device merges from, and a damaged one
/// costs far more than the check: each device downloads it, fails, and — for
/// a week, once — declines to push over it. `quick_check` reads every page
/// but skips index verification, which is the affordable version of "is this
/// a database" on a 40 MB file that has just been written and is still in the
/// page cache. A failure here is [`CatalogError::Corrupt`], the same thing a
/// receiving device would have said, so the sync reports it the same way.
fn verify_snapshot(snapshot: &Connection) -> Result<(), CatalogError> {
let verdict: String = snapshot.query_row("PRAGMA quick_check", [], |r| r.get(0))?;
if verdict == "ok" {
Ok(())
} else {
Err(CatalogError::Corrupt {
detail: format!("the snapshot for upload failed quick_check: {verdict}"),
})
}
}
/// Checkpoint, then copy the whole database to `dest`, and hand back the /// Checkpoint, then copy the whole database to `dest`, and hand back the
/// connection to the copy. /// connection to the copy.
/// ///
+118 -10
View File
@@ -121,16 +121,24 @@ impl NextcloudBackend {
body: Vec<u8>, body: Vec<u8>,
) -> Result<Validator, RemoteError> { ) -> Result<Validator, RemoteError> {
let total = body.len() as u64; let total = body.len() as u64;
// Named from the destination so a resumed or abandoned upload is // TRACES: FR-NC-9
// identifiable, and so two uploads cannot collide in one directory. // Named from the destination, so an abandoned upload is identifiable,
let token: String = path // **and from a nonce, so no two uploads ever share a directory.**
.as_str() //
.bytes() // The name used to be the destination alone, on the reasoning that two
.map(|b| match b { // *files* could then not collide. Two *devices* uploading the same file
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' => (b as char).to_string(), // could, and did: both wrote `00001`…`00009` into one directory, and
_ => "-".to_string(), // whichever `MOVE`d first assembled a mix of the two — a catalog of
}) // exactly the right size whose pages came from two different
.collect(); // databases. SQLite called it malformed, every client then declined to
// overwrite it, and collections stopped syncing on all of them for a
// week. An upload that died on a phone's link left its chunks there
// for the next device to assemble in, by the same mechanism.
let token = format!(
"{}-{}",
sanitise_for_upload_dir(path.as_str()),
upload_nonce()
);
let dir = format!( let dir = format!(
"{}/remote.php/dav/uploads/{}/{token}", "{}/remote.php/dav/uploads/{}/{token}",
self.server, self.login self.server, self.login
@@ -138,6 +146,27 @@ impl NextcloudBackend {
self.mkcol_url(&dir).await?; self.mkcol_url(&dir).await?;
// Whatever happens below, the directory does not outlive the attempt.
// With a unique name a leftover is only quota rather than corruption,
// but a phone that abandons uploads all day would still leave dozens
// of 5 MB chunks behind, and the server only sweeps them eventually.
let result = self.put_chunks_and_assemble(path, &dir, body, total).await;
if result.is_err() {
self.discard_upload_dir(&dir).await;
}
result
}
/// The transfer half of [`put_chunked`](Self::put_chunked): the chunks,
/// then the `MOVE` that assembles them. Split out so that a failure at any
/// point returns to one place that cleans up.
async fn put_chunks_and_assemble(
&self,
path: &RemotePath,
dir: &str,
body: Vec<u8>,
total: u64,
) -> Result<Validator, RemoteError> {
// Chunks are numbered from 1 and must sort correctly as strings, which // Chunks are numbered from 1 and must sort correctly as strings, which
// is why they are zero-padded rather than bare integers. // is why they are zero-padded rather than bare integers.
let chunk_size = CHUNKS.min_chunk as usize; let chunk_size = CHUNKS.min_chunk as usize;
@@ -191,6 +220,27 @@ impl NextcloudBackend {
self.dir_validator(path).await self.dir_validator(path).await
} }
/// Best-effort `DELETE` of an upload directory whose transfer failed.
///
/// Errors are logged and dropped: this runs on the way out of a failure,
/// and the failure is what the caller needs to hear about. A connection
/// that has just died will refuse this too, and that is fine — the
/// server sweeps abandoned upload directories on its own; this only
/// spares it the wait when the link is still up.
async fn discard_upload_dir(&self, dir: &str) {
let attempt = self
.client
.delete(dir)
.basic_auth(&self.login, Some(&self.password))
.send()
.await;
match attempt {
Ok(resp) if resp.status().is_success() || resp.status() == 404 => {}
Ok(resp) => log::debug!("leaving abandoned upload {dir}: {}", resp.status()),
Err(e) => log::debug!("leaving abandoned upload {dir}: {e}"),
}
}
/// `MKCOL` at an absolute URL, treating "already there" as success. /// `MKCOL` at an absolute URL, treating "already there" as success.
async fn mkcol_url(&self, url: &str) -> Result<(), RemoteError> { async fn mkcol_url(&self, url: &str) -> Result<(), RemoteError> {
let resp = self let resp = self
@@ -725,6 +775,37 @@ fn map_send_error(e: reqwest::Error) -> RemoteError {
} }
} }
/// The destination path reduced to characters every WebDAV server accepts in
/// an upload directory name.
fn sanitise_for_upload_dir(path: &str) -> String {
path.bytes()
.map(|b| match b {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' => (b as char).to_string(),
_ => "-".to_string(),
})
.collect()
}
/// A token no other upload — on this device or any other — will produce.
///
/// Nanoseconds since the epoch, the process id, and a counter, mixed rather
/// than concatenated so the name stays short. Two devices would have to start
/// an upload in the same nanosecond from the same pid to collide, and the
/// counter separates two uploads this process starts in one tick. No
/// randomness crate is pulled in for this: unique is the requirement, not
/// unguessable.
fn upload_nonce() -> String {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(0);
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0);
let pid = u64::from(std::process::id());
let n = COUNTER.fetch_add(1, Ordering::Relaxed);
format!("{:016x}", nanos ^ (pid << 40) ^ n.rotate_left(20))
}
/// Translate an HTTP status into a typed error. /// Translate an HTTP status into a typed error.
fn map_status(status: reqwest::StatusCode, what: &str) -> Result<(), RemoteError> { fn map_status(status: reqwest::StatusCode, what: &str) -> Result<(), RemoteError> {
match status.as_u16() { match status.as_u16() {
@@ -765,6 +846,33 @@ fn encode_path(path: &str) -> String {
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
#[test]
fn two_uploads_of_one_destination_never_share_a_directory() {
// The collision that assembled two devices' chunks into one file.
// Same path, back to back, same process: still two names.
let a = format!(
"{}-{}",
super::sanitise_for_upload_dir("PhotosRaw/.darkroom-derived/catalog.sqlite"),
super::upload_nonce()
);
let b = format!(
"{}-{}",
super::sanitise_for_upload_dir("PhotosRaw/.darkroom-derived/catalog.sqlite"),
super::upload_nonce()
);
assert_ne!(a, b);
assert!(a.starts_with("PhotosRaw--darkroom-derived-catalog-sqlite-"));
}
#[test]
fn upload_directory_names_are_plain() {
assert_eq!(
super::sanitise_for_upload_dir("Photos Raw/été/x.sqlite"),
"Photos-Raw---t---x-sqlite"
);
assert!(super::upload_nonce().bytes().all(|b| b.is_ascii_hexdigit()));
}
use super::*; use super::*;
fn backend() -> NextcloudBackend { fn backend() -> NextcloudBackend {
+24 -24
View File
File diff suppressed because one or more lines are too long
+557 -35
View File
@@ -52,6 +52,13 @@ pub struct SyncReport {
pub thumbnails_adopted: usize, pub thumbnails_adopted: usize,
pub catalog_uploaded: bool, pub catalog_uploaded: bool,
pub catalog_merged: 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, pub collections_gained: usize,
/// Membership rows the merge brought in (FR-CAT-7). /// Membership rows the merge brought in (FR-CAT-7).
/// ///
@@ -88,6 +95,7 @@ impl SyncReport {
|| self.shards_downloaded > 0 || self.shards_downloaded > 0
|| self.catalog_uploaded || self.catalog_uploaded
|| self.catalog_merged || self.catalog_merged
|| self.catalog_replaced
|| self.face_shards_uploaded > 0 || self.face_shards_uploaded > 0
|| self.face_shards_downloaded > 0 || self.face_shards_downloaded > 0
|| self.place_adopted || self.place_adopted
@@ -583,11 +591,28 @@ fn legacy_upload_of_ours(store: &ThumbStore, id: u32, remote_size: u64) -> bool
.unwrap_or(false) .unwrap_or(false)
} }
/// Exchange the catalog, for its collections. /// Exchange the catalog, for its collections and its people.
/// ///
/// Only collections merge — see [`dr_catalog::sync`]. The rest of a catalog /// Only collections, keywords and people merge — see [`dr_catalog::sync`]. The
/// describes local state (folder ETags, cache paths, job rows) and importing /// rest of a catalog describes local state (folder ETags, cache paths, job
/// another device's version would be actively wrong. /// rows) and importing another device's version would be actively wrong.
///
/// # Generations
///
/// TRACES: FR-NC-9 | NFR-R2
/// The server keeps the current copy and the [`GENERATIONS`] before it:
/// `catalog.sqlite`, then `catalog.1.sqlite` (the one it replaced), `.2`,
/// `.3`. A push uploads to a temporary name, rotates, and moves the upload into
/// place — so at no moment is there no current copy, and a copy that turns out
/// damaged has the one before it to fall back on. Rotation is server-side
/// renames; the only transfer is the upload itself.
///
/// This is what a damaged copy used to lack. With one copy and nothing behind
/// it, "the current file will not open" left two answers, both bad: refuse for
/// ever, or overwrite with ours and lose whatever another device had added
/// since. Now it has a third — merge from the newest readable generation, which
/// loses nothing — and the damaged file itself is kept as `.1` by the ordinary
/// rotation rather than by a separate upload.
async fn sync_catalog( async fn sync_catalog(
backend: &dyn RemoteBackend, backend: &dyn RemoteBackend,
base: &RemotePath, base: &RemotePath,
@@ -595,7 +620,7 @@ async fn sync_catalog(
scratch: &Path, scratch: &Path,
report: &mut SyncReport, report: &mut SyncReport,
) -> Result<(), String> { ) -> Result<(), String> {
let remote_name = "catalog.sqlite"; let remote_name = CATALOG_NAME;
let target = RemotePath::new(format!("{}/{remote_name}", base.as_str())); let target = RemotePath::new(format!("{}/{remote_name}", base.as_str()));
// ---- take theirs first ----------------------------------------------- // ---- take theirs first -----------------------------------------------
@@ -627,30 +652,48 @@ async fn sync_catalog(
}; };
if let Some(bytes) = theirs { if let Some(bytes) = theirs {
let downloaded = scratch.join("catalog-remote.sqlite"); let catalog = match dr_catalog::Catalog::open(catalog_path) {
if std::fs::write(&downloaded, &bytes).is_ok() { Ok(c) => c,
match dr_catalog::Catalog::open(catalog_path) { Err(e) => {
Ok(catalog) => match catalog.merge_remote_catalog(&downloaded) { log::warn!("not pushing the catalog: opening ours to merge: {e}");
Ok(merge) => { return Ok(());
report.catalog_merged = true; }
report.collections_gained = merge.inserted + merge.updated; };
report.members_gained = merge.members_added; match merge_downloaded(&catalog, scratch, &bytes, report) {
} Ok(()) => {}
// Unreadable is not the same as absent: it may be a newer // TRACES: FR-NC-9
// format, or a torn upload. Ours must not go over it. // Damaged beyond reading, and the whole file arrived — so it is
Err(e) => { // not a short download, and no device will read it either. Fall
log::warn!("not pushing the catalog: merging the server's copy: {e}"); // back to the generation before it; only with none readable is
let _ = std::fs::remove_file(&downloaded); // ours the whole truth. See the header for why refusing here for
return Ok(()); // ever was the wrong answer.
} Err(dr_catalog::CatalogError::Corrupt { detail }) => {
}, if !arrived_whole(backend, base, remote_name, bytes.len(), &detail).await {
Err(e) => {
log::warn!("not pushing the catalog: opening ours to merge: {e}");
let _ = std::fs::remove_file(&downloaded);
return Ok(()); return Ok(());
} }
match merge_from_generations(backend, base, &catalog, scratch, report).await {
Some(n) => log::warn!(
"the catalog on the server is damaged ({detail}) and all {} bytes of \
it arrived, so no device can read it; merged from the generation \
before it ({}) instead, and replacing it",
bytes.len(),
generation_name(n)
),
None => log::warn!(
"the catalog on the server is damaged ({detail}) and all {} bytes of \
it arrived, so no device can read it, and no earlier generation is \
readable either; replacing it with this device's copy",
bytes.len()
),
}
report.catalog_replaced = true;
}
// Unreadable is not the same as absent: it may be a newer
// format. Ours must not go over it.
Err(e) => {
log::warn!("not pushing the catalog: merging the server's copy: {e}");
return Ok(());
} }
let _ = std::fs::remove_file(&downloaded);
} }
} }
@@ -666,15 +709,177 @@ async fn sync_catalog(
.map_err(|e| e.to_string())?; .map_err(|e| e.to_string())?;
let bytes = std::fs::read(&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); let _ = std::fs::remove_file(&snapshot);
// To a temporary name first. The upload is the only step that can fail
// half-way, and a half-uploaded *current* copy is exactly the damaged file
// this whole scheme exists to survive. Under its own name a failure leaves
// the current copy untouched and costs one stray file, retried next pass.
let staging = RemotePath::new(format!("{}/{UPLOAD_NAME}", base.as_str()));
let sent = bytes.len() as u64;
if let Err(e) = backend.put(&staging, bytes, None).await {
log::warn!("uploading catalog: {e}");
return Ok(());
}
// TRACES: NFR-R2
// Confirm the server holds what was sent before it becomes the copy every
// other device reads. A chunked upload is assembled server-side, and an
// assembly that went wrong is a file of plausible size that no device can
// open — the one failure the generations exist to survive, and cheaper
// to catch here, on the device that caused it, than on every other one
// after. One listing; the size is what the server can vouch for without
// reading the file back.
match remote_size(backend, base, UPLOAD_NAME).await {
Some(held) if held == sent => {}
Some(held) => {
log::warn!(
"not replacing the catalog: sent {sent} bytes but the server holds {held}; \
the upload is discarded and retried next pass"
);
let _ = backend.delete(&RemoteId::Path(staging), None).await;
return Ok(());
}
None => {
log::warn!(
"not replacing the catalog: the server would not confirm the upload's size; \
it is discarded and retried next pass"
);
let _ = backend.delete(&RemoteId::Path(staging), None).await;
return Ok(());
}
}
// Rotate, then move the upload into place. Every step here is a rename on
// the server, and every destination is empty by the time it is written to
// — a `move_to` will not overwrite, by design — so a failure at any point
// leaves a gap in the generations and never a missing current copy.
if let Err(e) = rotate_generations(backend, base).await {
log::warn!("not replacing the catalog: rotating the earlier copies: {e}");
return Ok(());
}
match backend.move_to(&RemoteId::Path(staging), &target).await {
Ok(()) => report.catalog_uploaded = true,
Err(e) => log::warn!("moving the uploaded catalog into place: {e}"),
}
Ok(()) Ok(())
} }
/// The current copy's name on the server.
const CATALOG_NAME: &str = "catalog.sqlite";
/// Where a push lands before it is rotated into place.
const UPLOAD_NAME: &str = "catalog.upload.sqlite";
/// How many earlier copies the server keeps behind the current one.
///
/// Three, because what they are for is surviving one damaged push and the one
/// or two syncs it may take for a device to notice. More would cost nothing in
/// transfer — rotation is renames — but each is a 40 MB file on the account's
/// quota, and the local backups (NFR-R2) are the long-term store.
const GENERATIONS: usize = 3;
/// `catalog.N.sqlite` for `1 <= N <= GENERATIONS`; `.1` is the newest.
fn generation_name(n: usize) -> String {
format!("catalog.{n}.sqlite")
}
/// Merge a downloaded catalog into ours, through a file in scratch.
///
/// Attaching needs a path, and the download is bytes. Written and removed here
/// so the callers — the current copy, and each generation tried after it —
/// cannot disagree about cleanup.
fn merge_downloaded(
catalog: &dr_catalog::Catalog,
scratch: &Path,
bytes: &[u8],
report: &mut SyncReport,
) -> Result<(), dr_catalog::CatalogError> {
let downloaded = scratch.join("catalog-remote.sqlite");
std::fs::write(&downloaded, bytes).map_err(|e| dr_catalog::CatalogError::Io(e.to_string()))?;
let result = catalog.merge_remote_catalog(&downloaded);
let _ = std::fs::remove_file(&downloaded);
let merge = result?;
report.catalog_merged = true;
report.collections_gained += merge.inserted + merge.updated;
report.members_gained += merge.members_added;
Ok(())
}
/// TRACES: FR-NC-9
/// Merge from the newest generation that reads, when the current copy will
/// not. Returns which one, or `None` when none of them does.
///
/// Newest first, and the first readable one wins: a generation is a complete
/// snapshot, so an older one adds nothing a newer one lacks. A generation that
/// is absent, damaged, or from a newer schema is skipped the same way — none of
/// those is a reason to stop looking further back.
async fn merge_from_generations(
backend: &dyn RemoteBackend,
base: &RemotePath,
catalog: &dr_catalog::Catalog,
scratch: &Path,
report: &mut SyncReport,
) -> Option<usize> {
for n in 1..=GENERATIONS {
let path = RemotePath::new(format!("{}/{}", base.as_str(), generation_name(n)));
let bytes = match read_derived(backend, &path).await {
Ok(b) => b,
Err(RemoteError::NotFound(_)) => continue,
Err(e) => {
log::debug!("skipping {}: {e}", generation_name(n));
continue;
}
};
match merge_downloaded(catalog, scratch, &bytes, report) {
Ok(()) => return Some(n),
Err(e) => log::debug!("skipping {}: {e}", generation_name(n)),
}
}
None
}
/// Make room for a new current copy: drop the oldest generation and shift the
/// rest back by one, ending with the current copy as `.1`.
///
/// Oldest first, so that each destination is empty when it is moved into —
/// `move_to` refuses to overwrite, and rightly. A name that is not there is
/// not an error at any step: a library that has synced twice has no `.3` yet.
async fn rotate_generations(backend: &dyn RemoteBackend, base: &RemotePath) -> Result<(), String> {
let at = |name: String| RemotePath::new(format!("{}/{name}", base.as_str()));
match backend
.delete(&RemoteId::Path(at(generation_name(GENERATIONS))), None)
.await
{
Ok(()) | Err(RemoteError::NotFound(_)) => {}
Err(e) => return Err(format!("dropping {}: {e}", generation_name(GENERATIONS))),
}
for n in (1..GENERATIONS).rev() {
match backend
.move_to(
&RemoteId::Path(at(generation_name(n))),
&at(generation_name(n + 1)),
)
.await
{
Ok(()) | Err(RemoteError::NotFound(_)) => {}
Err(e) => return Err(format!("moving {} back: {e}", generation_name(n))),
}
}
match backend
.move_to(
&RemoteId::Path(at(CATALOG_NAME.into())),
&at(generation_name(1)),
)
.await
{
Ok(()) | Err(RemoteError::NotFound(_)) => Ok(()),
Err(e) => Err(format!("setting the current copy back: {e}")),
}
}
/// TRACES: FR-UI-8 /// TRACES: FR-UI-8
/// The place's name inside the derived folder. /// The place's name inside the derived folder.
/// ///
@@ -851,6 +1056,64 @@ async fn read_derived(
} }
} }
/// TRACES: FR-NC-9
/// Whether a download that will not open is the server's whole file.
///
/// A truncated download is also unreadable, and on a phone it is the far
/// likelier story — this library's logs are full of aborted bodies and DNS
/// failures. Treating that as damage would let one bad connection discard a
/// catalog the server was holding perfectly well. So the size the server
/// advertises is compared against what actually arrived, and anything short,
/// or any size the listing cannot confirm, is the transport failure it is:
/// the caller must leave the server's copy alone.
async fn arrived_whole(
backend: &dyn RemoteBackend,
base: &RemotePath,
name: &str,
received: usize,
detail: &str,
) -> bool {
let Some(advertised) = remote_size(backend, base, name).await else {
log::warn!(
"not pushing the catalog: the copy on the server will not open ({detail}), \
but its size could not be confirmed, so it may simply have arrived short"
);
return false;
};
if advertised != received as u64 {
log::warn!(
"not pushing the catalog: the copy on the server will not open ({detail}), \
but only {received} of {advertised} bytes arrived — that is a truncated \
download, not a damaged file, so the server's copy is left alone"
);
return false;
}
true
}
/// The size the server says an entry in the derived folder has.
///
/// One listing of a folder that holds a handful of files, rather than a HEAD
/// the backend trait does not offer. `None` covers "the listing failed", "it is
/// not there", and "the size it reports means nothing" alike, and the caller
/// treats all of them as not knowing — which is the answer that declines to
/// overwrite.
///
/// A placeholder is excluded rather than trusted: `RemoteEntry::size` is
/// explicitly not meaningful when `materialised` is false — a suffix-mode stub
/// is one byte and carries no record of what it stands for — so comparing a
/// download against it would be comparing against nothing.
async fn remote_size(backend: &dyn RemoteBackend, base: &RemotePath, name: &str) -> Option<u64> {
backend
.list(base, None)
.await
.ok()?
.into_iter()
.find(|e| e.kind == dr_sync::EntryKind::File && e.path.name() == name)
.filter(|e| e.materialised)
.map(|e| e.size)
}
fn shard_name(client: &str, id: u32) -> String { fn shard_name(client: &str, id: u32) -> String {
format!("shard-{client}-{id:04}.sqlite") format!("shard-{client}-{id:04}.sqlite")
} }
@@ -979,6 +1242,19 @@ mod derived_guard_tests {
/// record went up rather than only that one did. /// record went up rather than only that one did.
last_put: Arc<std::sync::Mutex<Vec<u8>>>, last_put: Arc<std::sync::Mutex<Vec<u8>>>,
caps: dr_sync::Capabilities, caps: dr_sync::Capabilities,
/// The size `list` claims `catalog.sqlite` has, if it lists it at all.
/// This is what says whether a body that will not open is a damaged
/// file or merely a short download.
advertise: Option<u64>,
/// Bodies by file name, for tests that need the current copy and a
/// generation to differ. When empty every read serves `body`; when
/// not, a name absent here reads as `NotFound`.
bodies: std::collections::HashMap<String, Vec<u8>>,
/// Every `move_to`, as (from, to) names, in order.
moves: Arc<std::sync::Mutex<Vec<(String, String)>>>,
/// What `list` claims the staged upload's size is, when a test wants
/// the server to have assembled it wrongly. `None` reports the truth.
staged_size: Option<u64>,
} }
impl Fussy { impl Fussy {
@@ -1000,6 +1276,10 @@ mod derived_guard_tests {
puts: puts.clone(), puts: puts.clone(),
last_put: last_put.clone(), last_put: last_put.clone(),
caps: dr_sync::Capabilities::minimal(), caps: dr_sync::Capabilities::minimal(),
advertise: None,
bodies: Default::default(),
moves: Default::default(),
staged_size: None,
}, },
puts, puts,
last_put, last_put,
@@ -1017,10 +1297,36 @@ mod derived_guard_tests {
} }
async fn list( async fn list(
&self, &self,
_dir: &RemotePath, dir: &RemotePath,
_since: Option<&dr_sync::Validator>, _since: Option<&dr_sync::Validator>,
) -> Result<Vec<dr_sync::RemoteEntry>, RemoteError> { ) -> Result<Vec<dr_sync::RemoteEntry>, RemoteError> {
Ok(Vec::new()) let entry = |name: &str, size: u64| {
let path = RemotePath::new(format!("{}/{name}", dir.as_str()));
dr_sync::RemoteEntry {
id: RemoteId::Path(path.clone()),
path,
kind: dr_sync::EntryKind::File,
validator: dr_sync::Validator::new("v"),
size,
modified: None,
has_preview: false,
materialised: true,
}
};
let mut out = Vec::new();
if let Some(size) = self.advertise {
out.push(entry(CATALOG_NAME, size));
}
// The staged upload lists at the size of what was last put, as a
// server that assembled it correctly would report — unless a test
// says the assembly went wrong.
if self.puts.load(Ordering::SeqCst) > 0 {
let held = self
.staged_size
.unwrap_or(self.last_put.lock().unwrap().len() as u64);
out.push(entry(UPLOAD_NAME, held));
}
Ok(out)
} }
async fn dir_validator( async fn dir_validator(
&self, &self,
@@ -1036,7 +1342,7 @@ mod derived_guard_tests {
} }
async fn get( async fn get(
&self, &self,
_id: &RemoteId, id: &RemoteId,
_r: Option<std::ops::Range<u64>>, _r: Option<std::ops::Range<u64>>,
) -> Result<Vec<u8>, RemoteError> { ) -> Result<Vec<u8>, RemoteError> {
match &self.fail_with { match &self.fail_with {
@@ -1045,7 +1351,17 @@ mod derived_guard_tests {
Err(RemoteError::NotMaterialised(s.clone())) Err(RemoteError::NotMaterialised(s.clone()))
} }
Some(_) => Err(RemoteError::PermissionDenied), Some(_) => Err(RemoteError::PermissionDenied),
None => Ok(self.body.clone()), None if self.bodies.is_empty() => Ok(self.body.clone()),
None => {
let name = match id {
RemoteId::Path(p) => p.name().to_string(),
_ => String::new(),
};
self.bodies
.get(&name)
.cloned()
.ok_or(RemoteError::NotFound(name))
}
} }
} }
async fn put( async fn put(
@@ -1065,7 +1381,15 @@ mod derived_guard_tests {
) -> Result<(), RemoteError> { ) -> Result<(), RemoteError> {
Ok(()) Ok(())
} }
async fn move_to(&self, _f: &RemoteId, _t: &RemotePath) -> Result<(), RemoteError> { async fn move_to(&self, f: &RemoteId, t: &RemotePath) -> Result<(), RemoteError> {
let from = match f {
RemoteId::Path(p) => p.name().to_string(),
_ => String::new(),
};
self.moves
.lock()
.unwrap()
.push((from, t.name().to_string()));
Ok(()) Ok(())
} }
async fn create_dir(&self, _p: &RemotePath) -> Result<(), RemoteError> { async fn create_dir(&self, _p: &RemotePath) -> Result<(), RemoteError> {
@@ -1135,6 +1459,204 @@ mod derived_guard_tests {
assert!(report.catalog_uploaded); assert!(report.catalog_uploaded);
} }
// --- a damaged copy on the server (FR-NC-9) --------------------------
//
// The one exception to everything above. A file SQLite calls malformed is
// not a read that failed; it is a file no device will ever read again, and
// leaving it alone pins it there for every client at once.
/// The bytes of a catalog holding one collection, as a peer would upload.
fn a_catalog_with(collection: &str, name: &str) -> Vec<u8> {
let dir = std::env::temp_dir().join(format!("dr-generation-{name}"));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("catalog.sqlite");
let cat = dr_catalog::Catalog::open(&path).unwrap();
cat.connection()
.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES (?1, ?2, 0, 0, 1, 1)",
rusqlite::params![format!("uuid-{collection}"), collection],
)
.unwrap();
let snap = dir.join("snap.sqlite");
cat.snapshot_for_upload(&snap).unwrap();
let bytes = std::fs::read(&snap).unwrap();
let _ = std::fs::remove_dir_all(&dir);
bytes
}
/// A body that is definitely not a SQLite database as the current copy,
/// with a chosen advertised size and whatever generations the test wants.
async fn corrupt_remote(
body_len: usize,
advertised: Option<u64>,
generations: &[(usize, Vec<u8>)],
name: &str,
) -> (usize, SyncReport, Vec<(String, String)>) {
let (catalog_path, scratch) = fixture(name);
let (mut backend, puts, _) = Fussy::serving(None, Vec::new());
backend.advertise = advertised;
backend
.bodies
.insert(CATALOG_NAME.to_string(), vec![0xAB; body_len]);
for (n, bytes) in generations {
backend.bodies.insert(generation_name(*n), bytes.clone());
}
let moves = backend.moves.clone();
let mut report = SyncReport::default();
sync_catalog(
&backend,
&RemotePath::new(".darkroom-derived"),
&catalog_path,
&scratch,
&mut report,
)
.await
.unwrap();
let moves = moves.lock().unwrap().clone();
(puts.load(Ordering::SeqCst), report, moves)
}
#[tokio::test]
async fn a_damaged_catalog_that_arrived_whole_is_replaced() {
// The deadlock this exists to break: every device downloads the same
// unreadable file, every device declines to overwrite it, and
// collections and people stop crossing between devices for ever.
let (puts, report, moves) = corrupt_remote(64, Some(64), &[], "corrupt-whole").await;
assert_eq!(puts, 1, "ours goes up, to the staging name");
assert!(report.catalog_replaced, "and the report says what happened");
assert!(report.catalog_uploaded);
assert!(
!report.catalog_merged,
"there was nothing readable to merge"
);
// The damaged file is kept by the rotation, not thrown away.
assert!(moves.contains(&(CATALOG_NAME.into(), generation_name(1))));
assert_eq!(
moves.last().unwrap(),
&(UPLOAD_NAME.to_string(), CATALOG_NAME.to_string()),
"and the upload is moved into place last"
);
}
#[tokio::test]
async fn a_damaged_catalog_falls_back_to_the_generation_before_it() {
// What generations are for. The device that pushed the damaged copy
// may have been the only one holding some collection; the generation
// before it still has everything every device had agreed on.
let older = a_catalog_with("Iceland", "gen1");
let (puts, report, _) =
corrupt_remote(64, Some(64), &[(1, older)], "corrupt-with-gen").await;
assert!(report.catalog_merged, "merged from catalog.1.sqlite");
assert_eq!(report.collections_gained, 1, "and gained what it held");
assert!(report.catalog_replaced);
assert_eq!(puts, 1);
}
#[tokio::test]
async fn a_damaged_generation_is_skipped_for_the_one_behind_it() {
// Two bad pushes in a row must not be worse than one.
let older = a_catalog_with("Faroe", "gen2");
let (_, report, _) = corrupt_remote(
64,
Some(64),
&[(1, vec![0xCD; 64]), (2, older)],
"corrupt-two-deep",
)
.await;
assert!(report.catalog_merged);
assert_eq!(report.collections_gained, 1);
}
#[tokio::test]
async fn a_short_download_is_a_truncated_transfer_not_a_damaged_file() {
// The failure that matters most to get right. A phone on a flaky link
// aborts bodies constantly, and a partial download will not open
// either — treating that as damage would let one bad connection
// destroy a catalog the server was holding perfectly well.
let (puts, report, moves) = corrupt_remote(64, Some(4096), &[], "corrupt-short").await;
assert_eq!(puts, 0, "nothing is written over a copy that arrived short");
assert!(moves.is_empty(), "and nothing is rotated");
assert!(!report.catalog_replaced);
assert!(!report.catalog_uploaded);
}
#[tokio::test]
async fn a_size_the_server_will_not_confirm_leaves_the_copy_alone() {
// Not knowing is not the same as knowing it is whole. Without a size
// to compare against there is no way to tell damage from truncation,
// and the answer to "I cannot tell" has to stay "do not overwrite".
let (puts, report, _) = corrupt_remote(64, None, &[], "corrupt-unconfirmed").await;
assert_eq!(puts, 0);
assert!(!report.catalog_replaced);
}
#[tokio::test]
async fn a_push_rotates_oldest_first_and_lands_last() {
// The order is the safety: every destination is empty when it is
// moved into, so a failure at any step leaves a gap and never a
// missing current copy.
let (puts, report, moves) =
run_with_moves(Some(RemoteError::NotFound("nope".into())), "rotation").await;
assert_eq!(puts, 1);
assert!(report.catalog_uploaded);
assert_eq!(
moves,
vec![
(generation_name(2), generation_name(3)),
(generation_name(1), generation_name(2)),
(CATALOG_NAME.to_string(), generation_name(1)),
(UPLOAD_NAME.to_string(), CATALOG_NAME.to_string()),
]
);
}
#[tokio::test]
async fn an_upload_the_server_holds_at_the_wrong_size_is_not_rotated_in() {
// The assembly went wrong on the server. Rotating that into place would
// hand every other device the damaged file the generations exist to
// survive; discarding it costs one retry.
let (catalog_path, scratch) = fixture("wrong-size");
let (mut backend, puts) = Fussy::reading(Some(RemoteError::NotFound("nope".into())));
backend.staged_size = Some(7);
let moves = backend.moves.clone();
let mut report = SyncReport::default();
sync_catalog(
&backend,
&RemotePath::new(".darkroom-derived"),
&catalog_path,
&scratch,
&mut report,
)
.await
.unwrap();
assert_eq!(puts.load(Ordering::SeqCst), 1, "it was uploaded");
assert!(!report.catalog_uploaded, "but not accepted");
assert!(moves.lock().unwrap().is_empty(), "and nothing was rotated");
}
async fn run_with_moves(
fail_with: Option<RemoteError>,
name: &str,
) -> (usize, SyncReport, Vec<(String, String)>) {
let (catalog_path, scratch) = fixture(name);
let (backend, puts) = Fussy::reading(fail_with);
let moves = backend.moves.clone();
let mut report = SyncReport::default();
sync_catalog(
&backend,
&RemotePath::new(".darkroom-derived"),
&catalog_path,
&scratch,
&mut report,
)
.await
.unwrap();
let moves = moves.lock().unwrap().clone();
(puts.load(Ordering::SeqCst), report, moves)
}
// --- the place (FR-UI-8) --------------------------------------------- // --- the place (FR-UI-8) ---------------------------------------------
// //
// The same "do not write over what you could not read" rule as above, for a // The same "do not write over what you could not read" rule as above, for a
+44 -4
View File
@@ -4155,6 +4155,34 @@ fn apply_zoom(window: &AppWindow, ctl: &Rc<LibraryController>, delta: i32) {
refresh_timeline(window, catalog, ctl); refresh_timeline(window, catalog, ctl);
} }
/// TRACES: NFR-R2
/// Take the day's catalog backup if one is due, off the UI thread.
///
/// Decided cheaply first — [`dr_catalog::recovery::backup_due`] reads a
/// directory listing, not the catalog — so the ordinary case of "backed up
/// this morning already" costs no thread and no connection.
fn spawn_scheduled_backup(ctl: &Rc<LibraryController>) {
let Some((conn, _)) = ctl.session.borrow().clone() else {
return;
};
let catalog_path = library::catalog_path(&conn.account);
if !dr_catalog::recovery::backup_due(&catalog_path) {
return;
}
std::thread::spawn(move || {
let catalog = match dr_catalog::Catalog::open(&catalog_path) {
Ok(c) => c,
Err(e) => {
log::warn!("scheduled backup: opening the catalog: {e}");
return;
}
};
if let Err(e) = dr_catalog::recovery::backup_if_due(catalog.connection(), &catalog_path) {
log::warn!("scheduled backup: {e}");
}
});
}
/// Push shards and the catalog to the server, and take what it has. /// Push shards and the catalog to the server, and take what it has.
/// ///
/// Fired after the sweep completes, when there is a finished index worth /// Fired after the sweep completes, when there is a finished index worth
@@ -4286,10 +4314,14 @@ fn start_derived_sync(window: &AppWindow, ctl: &Rc<LibraryController>) {
report.shards_uploaded, report.shards_uploaded,
report.shards_downloaded, report.shards_downloaded,
report.thumbnails_adopted, report.thumbnails_adopted,
if report.catalog_uploaded { // Replacing a damaged copy is said apart from an
"pushed" // ordinary push: it is the one push that discarded
} else { // something, and a log line that called it "pushed"
"not pushed" // would hide the only moment worth going back to.
match (report.catalog_uploaded, report.catalog_replaced) {
(true, true) => "pushed over a damaged copy",
(true, false) => "pushed",
_ => "not pushed",
}, },
if report.collections_gained > 0 { if report.collections_gained > 0 {
format!(", {} collection(s) gained", report.collections_gained) format!(", {} collection(s) gained", report.collections_gained)
@@ -4441,6 +4473,14 @@ fn start_sweep(window: &AppWindow, ctl: &Rc<LibraryController>) {
refresh_timeline(&w, catalog, &ctl_cb); refresh_timeline(&w, catalog, &ctl_cb);
} }
} }
// TRACES: NFR-R2
// The scheduled backup, at the moment the catalog is
// quietest and a day's edits have just been folded in.
// On its own thread, with its own connection: a copy
// of a 130 MB file is a second or two the Slint loop
// must not spend, and it runs beside the sync below,
// which is only a reader of the same file.
spawn_scheduled_backup(&ctl_cb);
// Now that indexing is complete, hand the result to the // Now that indexing is complete, hand the result to the
// server so a second device inherits it rather than // server so a second device inherits it rather than
// repeating hours of range fetches. // repeating hours of range fetches.