Two checks around the upload, both cheap next to what they prevent. Before: the snapshot is quick_checked before it leaves. It is the copy every other device merges from, and a damaged one costs each of them a download, a failed merge and a refusal to push. After: the staged upload's size on the server is compared to the bytes sent before it is rotated into place. A chunked upload is assembled server-side, and an assembly that goes wrong is a file of plausible size no device can open — caught here, on the device that caused it, for one listing; otherwise on every other device, after the fact. A mismatch, or a size the server will not confirm, discards the upload and leaves the current copy and its generations untouched.
384 lines
14 KiB
Rust
384 lines
14 KiB
Rust
//! TRACES: FR-CAT-7 | FR-NC-9 | NFR-R1
|
|
//! Preparing the catalog file for upload, and taking in a remote one.
|
|
//!
|
|
//! # The hazard this module exists to handle
|
|
//!
|
|
//! A WAL-mode SQLite database is not one file. Committed transactions can live
|
|
//! in `catalog.sqlite-wal` with the main file lagging behind, so copying
|
|
//! `catalog.sqlite` alone uploads a **torn snapshot**: internally consistent as
|
|
//! of some older point, missing everything since. Worse, a naive copy taken
|
|
//! while a writer is mid-transaction can be structurally corrupt.
|
|
//!
|
|
//! So an upload never copies the live file. It runs a TRUNCATE checkpoint to
|
|
//! fold the WAL back into the main file, then uses SQLite's own backup API to
|
|
//! take a consistent snapshot — which serialises correctly against concurrent
|
|
//! writers rather than racing them.
|
|
//!
|
|
//! # What is actually synced
|
|
//!
|
|
//! Only the *user's judgements about their library* merge: collections, and the
|
|
//! keyword vocabulary with its assignments (see [`crate::merge`]). The rest of
|
|
//! the catalog is a *local index* of *local* storage — folder mtimes, cache
|
|
//! paths, job rows — and copying another device's version of those in would be
|
|
//! actively wrong. The remote file is read for those two and then discarded.
|
|
//!
|
|
//! This is why the catalog remains disposable in the ARCH §6.12 sense: nothing
|
|
//! here makes the local database authoritative for anything a rebuild could
|
|
//! not recover.
|
|
|
|
use std::path::{Path, PathBuf};
|
|
|
|
use rusqlite::Connection;
|
|
|
|
use crate::error::CatalogError;
|
|
use crate::merge::{self, MergeReport};
|
|
|
|
/// Schema name the downloaded remote catalog is attached under.
|
|
const REMOTE_SCHEMA: &str = "remote_cat";
|
|
|
|
/// Fold the WAL into the main database file.
|
|
///
|
|
/// TRUNCATE rather than PASSIVE: passive checkpointing gives up when a reader
|
|
/// holds the WAL open, which would leave recent commits out of the snapshot
|
|
/// without saying so.
|
|
pub fn checkpoint(conn: &Connection) -> Result<(), CatalogError> {
|
|
conn.pragma_update(None, "wal_checkpoint", "TRUNCATE")?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Write a consistent snapshot of the catalog to `dest`, ready to upload.
|
|
///
|
|
/// Uses the backup API rather than a filesystem copy so the snapshot is
|
|
/// coherent even with writers active. Callers should still prefer a quiet
|
|
/// moment — this competes with background jobs for the write lock.
|
|
pub fn snapshot_for_upload(conn: &Connection, dest: &Path) -> Result<(), CatalogError> {
|
|
let out = copy_to(conn, dest)?;
|
|
strip_face_crops(&out)?;
|
|
verify_snapshot(&out)?;
|
|
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
|
|
/// connection to the copy.
|
|
///
|
|
/// Split out from [`snapshot_for_upload`] because [`crate::recovery`] wants
|
|
/// exactly this and none of what follows it there: an NFR-R2 backup is the
|
|
/// file the user may have to *live on*, so it keeps the face crops that an
|
|
/// upload strips. Sharing the copy rather than reimplementing it is what keeps
|
|
/// the WAL discipline in one place — a backup taken with `fs::copy` would be
|
|
/// the torn snapshot this module's header exists to warn about.
|
|
pub(crate) fn copy_to(conn: &Connection, dest: &Path) -> Result<Connection, CatalogError> {
|
|
checkpoint(conn)?;
|
|
|
|
let mut out = Connection::open(dest)?;
|
|
let backup = rusqlite::backup::Backup::new(conn, &mut out)?;
|
|
// SQLite's own "copy everything" sentinel is -1, but rusqlite asserts a
|
|
// positive page count, so ask for more pages than a catalog will ever
|
|
// have. The effect is the same: one step, no interleaved writers, no
|
|
// progress callback. A 50k-image catalog is tens of megabytes.
|
|
backup.run_to_completion(i32::MAX, std::time::Duration::ZERO, None)?;
|
|
drop(backup);
|
|
|
|
Ok(out)
|
|
}
|
|
|
|
/// Drop the stored face crops from a snapshot before it is uploaded.
|
|
///
|
|
/// The snapshot is the *whole catalog*, uploaded on every sync and downloaded
|
|
/// by every device. Face crops are a few KB each and a fully indexed library
|
|
/// holds tens of thousands of them, so leaving them in would put tens of MB on
|
|
/// every round trip — the exact cost `face_shard`'s 25 MB cap exists to bound,
|
|
/// and the reason the bulk per-face data lives in shards in the first place.
|
|
///
|
|
/// Crops are not lost by this: they travel in the face shards
|
|
/// ([`crate::face_shard::export_to_shards`]), which are written once and
|
|
/// downloaded once. Nothing reads a crop out of a merged remote catalog —
|
|
/// [`merge_all`] touches collections and keywords only — so removing them here
|
|
/// costs a receiving device nothing it would otherwise have had.
|
|
///
|
|
/// `VACUUM` afterwards because SQLite does not return freed pages to the file
|
|
/// on its own, and an upload sized by the file rather than by its contents
|
|
/// would keep paying for bytes that are no longer there.
|
|
fn strip_face_crops(snapshot: &Connection) -> Result<(), CatalogError> {
|
|
// A catalog older than the crop column is a legitimate input here — a
|
|
// snapshot taken mid-migration, or a test fixture built from an earlier
|
|
// schema — so an absent column is nothing to fail over.
|
|
let has_crop = snapshot
|
|
.prepare("SELECT crop FROM faces LIMIT 1")
|
|
.map(|_| true)
|
|
.unwrap_or(false);
|
|
if !has_crop {
|
|
return Ok(());
|
|
}
|
|
snapshot.execute("UPDATE faces SET crop = NULL WHERE crop IS NOT NULL", [])?;
|
|
snapshot.execute_batch("VACUUM")?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Whether a downloaded remote catalog is worth merging.
|
|
///
|
|
/// Cheap guard before attaching: a remote written by a newer build may contain
|
|
/// tables and columns this one cannot read, and attempting the merge would
|
|
/// fail mid-transaction rather than declining cleanly.
|
|
pub fn remote_is_mergeable(remote: &Path) -> Result<bool, CatalogError> {
|
|
let conn = Connection::open_with_flags(
|
|
remote,
|
|
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
|
|
)?;
|
|
let v: i64 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?;
|
|
Ok(v <= crate::schema::SCHEMA_VERSION)
|
|
}
|
|
|
|
/// Attach a downloaded remote catalog, merge its collections, detach.
|
|
///
|
|
/// The remote file is opened **read-only** — this device never writes to
|
|
/// another device's catalog, it only reads collections out of it.
|
|
pub fn merge_remote(conn: &Connection, remote: &Path) -> Result<MergeReport, CatalogError> {
|
|
if !remote_is_mergeable(remote)? {
|
|
return Err(CatalogError::SchemaTooNew {
|
|
found: -1,
|
|
supported: crate::schema::SCHEMA_VERSION,
|
|
});
|
|
}
|
|
|
|
// Path binds as a parameter; ATTACH accepts one, so a path containing a
|
|
// quote cannot break out into SQL.
|
|
conn.execute(
|
|
&format!("ATTACH DATABASE ?1 AS {REMOTE_SCHEMA}"),
|
|
[remote.to_string_lossy().as_ref()],
|
|
)?;
|
|
|
|
let result = merge::merge_all(conn);
|
|
|
|
// Detach even if the merge failed, or the next attempt errors with
|
|
// "database remote_cat is already in use".
|
|
let detach = conn.execute(&format!("DETACH DATABASE {REMOTE_SCHEMA}"), []);
|
|
if let Err(e) = detach {
|
|
log::warn!("failed to detach remote catalog: {e}");
|
|
}
|
|
|
|
result
|
|
}
|
|
|
|
/// Where the catalog snapshot and the downloaded remote live.
|
|
///
|
|
/// Both are transient working files, not the catalog itself, so they belong in
|
|
/// the cache directory rather than beside the live database.
|
|
#[derive(Debug, Clone)]
|
|
pub struct SyncPaths {
|
|
pub upload_snapshot: PathBuf,
|
|
pub downloaded_remote: PathBuf,
|
|
}
|
|
|
|
impl SyncPaths {
|
|
pub fn in_dir(cache_dir: &Path) -> Self {
|
|
SyncPaths {
|
|
upload_snapshot: cache_dir.join("catalog-upload.sqlite"),
|
|
downloaded_remote: cache_dir.join("catalog-remote.sqlite"),
|
|
}
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::schema;
|
|
|
|
fn seeded(path: &Path) -> Connection {
|
|
let c = Connection::open(path).unwrap();
|
|
schema::configure(&c).unwrap();
|
|
schema::migrate(&c).unwrap();
|
|
c
|
|
}
|
|
|
|
#[test]
|
|
fn snapshot_captures_committed_data() {
|
|
let dir = tempdir();
|
|
let live = dir.join("catalog.sqlite");
|
|
let snap = dir.join("snap.sqlite");
|
|
|
|
let c = seeded(&live);
|
|
c.execute(
|
|
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
|
|
VALUES ('u1', 'Iceland', 0, 0, 1, 1)",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
|
|
snapshot_for_upload(&c, &snap).unwrap();
|
|
|
|
// The snapshot must hold the row even though it was written after the
|
|
// database was created — the torn-file failure this guards against.
|
|
let s = Connection::open(&snap).unwrap();
|
|
let name: String = s
|
|
.query_row("SELECT name FROM collections", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(name, "Iceland");
|
|
}
|
|
|
|
#[test]
|
|
fn a_remote_from_a_newer_build_is_declined_not_attempted() {
|
|
let dir = tempdir();
|
|
let remote = dir.join("remote.sqlite");
|
|
let r = seeded(&remote);
|
|
r.pragma_update(None, "user_version", schema::SCHEMA_VERSION + 1)
|
|
.unwrap();
|
|
drop(r);
|
|
|
|
assert!(!remote_is_mergeable(&remote).unwrap());
|
|
|
|
let local = seeded(&dir.join("local.sqlite"));
|
|
assert!(matches!(
|
|
merge_remote(&local, &remote),
|
|
Err(CatalogError::SchemaTooNew { .. })
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn merge_remote_round_trips_a_collection() {
|
|
let dir = tempdir();
|
|
let remote_path = dir.join("remote.sqlite");
|
|
{
|
|
let r = seeded(&remote_path);
|
|
r.execute(
|
|
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
|
|
VALUES ('u-remote', 'Portugal', 0, 0, 1, 1)",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
checkpoint(&r).unwrap();
|
|
}
|
|
|
|
let local = seeded(&dir.join("local.sqlite"));
|
|
local
|
|
.execute(
|
|
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
|
|
VALUES ('u-local', 'Iceland', 0, 0, 1, 1)",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
|
|
let report = merge_remote(&local, &remote_path).unwrap();
|
|
assert_eq!(report.inserted, 1);
|
|
|
|
let n: i64 = local
|
|
.query_row("SELECT count(*) FROM collections", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(n, 2);
|
|
}
|
|
|
|
#[test]
|
|
fn the_remote_can_be_merged_twice_without_attach_conflict() {
|
|
// Detach must happen even on the failure path, or the second attempt
|
|
// errors with "database remote_cat is already in use".
|
|
let dir = tempdir();
|
|
let remote_path = dir.join("remote.sqlite");
|
|
{
|
|
let r = seeded(&remote_path);
|
|
r.execute(
|
|
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
|
|
VALUES ('u-remote', 'Portugal', 0, 0, 1, 1)",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
checkpoint(&r).unwrap();
|
|
}
|
|
let local = seeded(&dir.join("local.sqlite"));
|
|
|
|
merge_remote(&local, &remote_path).unwrap();
|
|
let second = merge_remote(&local, &remote_path).unwrap();
|
|
assert!(!second.local_changed());
|
|
}
|
|
|
|
/// A scratch directory that cleans up with the test.
|
|
fn tempdir() -> PathBuf {
|
|
let base = std::env::temp_dir().join(format!(
|
|
"dr-catalog-test-{}-{:?}",
|
|
std::process::id(),
|
|
std::thread::current().id()
|
|
));
|
|
let _ = std::fs::remove_dir_all(&base);
|
|
std::fs::create_dir_all(&base).unwrap();
|
|
base
|
|
}
|
|
|
|
/// The whole reason crops live in the shards: a snapshot is uploaded whole,
|
|
/// on every sync, to every device.
|
|
#[test]
|
|
fn the_snapshot_carries_no_face_crops() {
|
|
let dir = tempdir();
|
|
let live = dir.join("catalog.sqlite");
|
|
let snap = dir.join("snap.sqlite");
|
|
|
|
let c = seeded(&live);
|
|
c.execute(
|
|
"INSERT OR IGNORE INTO roots(id, kind, label) VALUES (1, 'local', 'lib')",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
c.execute(
|
|
"INSERT INTO images(id, root_id, source_ref, added_at) VALUES (1, 1, 'a.CR3', 0)",
|
|
[],
|
|
)
|
|
.unwrap();
|
|
c.execute(
|
|
"INSERT INTO faces
|
|
(image_id, x, y, w, h, landmarks, detector_confidence, embedding,
|
|
crop_px, model_id, detected_at, crop)
|
|
VALUES (1, 0.1, 0.1, 0.2, 0.2, X'00', 0.9, X'00', 180.0, 'm', 0, ?1)",
|
|
[vec![7u8; 4096]],
|
|
)
|
|
.unwrap();
|
|
|
|
snapshot_for_upload(&c, &snap).unwrap();
|
|
|
|
let out = Connection::open(&snap).unwrap();
|
|
let crops: i64 = out
|
|
.query_row(
|
|
"SELECT COUNT(*) FROM faces WHERE crop IS NOT NULL",
|
|
[],
|
|
|r| r.get(0),
|
|
)
|
|
.unwrap();
|
|
assert_eq!(crops, 0, "the snapshot still carries face crops");
|
|
|
|
// The face itself must still be there — only the pixels are dropped.
|
|
let faces: i64 = out
|
|
.query_row("SELECT COUNT(*) FROM faces", [], |r| r.get(0))
|
|
.unwrap();
|
|
assert_eq!(faces, 1);
|
|
|
|
// And the local catalog keeps its crop: this strips the copy, never
|
|
// the original.
|
|
let kept: i64 = c
|
|
.query_row(
|
|
"SELECT COUNT(*) FROM faces WHERE crop IS NOT NULL",
|
|
[],
|
|
|r| r.get(0),
|
|
)
|
|
.unwrap();
|
|
assert_eq!(kept, 1, "stripping the snapshot damaged the live catalog");
|
|
}
|
|
}
|