Files
DarkRoom/core/dr-catalog/src/sync.rs
T
dtourolle 78df4211b0 Run the people deduplication after every catalog merge (#78)
A merge is where two devices' people meet: the same name typed on each,
or a redirect one of them made. So dedup_people::run now follows every
successful sync::merge_remote. It runs on the sync worker, never the
UI thread, before the snapshot is pushed, so what it folds reaches the
server on the same pass.

It runs in its own transaction, and a failure is logged, not
returned. What the merge took is committed and valid either way, and
the next pass tries again. Once a catalog is clean it costs 8-15 ms on
the reference library (19k faces, 26k people rows). On copies of the
two real catalogs, a round trip with a peer running the previous merge
converges on 75 listed named people on both sides and stays there
over a second round. merge_remote with the job takes 75-90 ms there.
2026-09-26 17:50:38 -04:00

832 lines
32 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,
/// without the face crops.
///
/// Built rather than copied. The snapshot is the *whole catalog* bar the
/// crops, uploaded on every sync and downloaded by every device. The crops are
/// most of the file (96 MB of a 158 MB reference catalog), and copying them in
/// only to delete them was most of the cost. A backup-API copy followed by
/// `UPDATE faces SET crop = NULL` and `VACUUM` wrote the file roughly three
/// times over to produce 50 MB (#71). So this creates the schema in an empty
/// file and copies every table into it with `crop` left NULL. That is one
/// pass, with nothing written that is not uploaded.
///
/// Consistency comes from doing the whole copy inside one transaction on the
/// snapshot's connection, which holds a single read snapshot of the source for
/// its duration. A writer committing meanwhile lands in the source's WAL and is
/// simply not seen, the same serialisation the backup API gave.
///
/// 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. The
/// merge reads a remote face's box and model to match it to a local one, and
/// no more. So leaving them out costs a receiving device nothing it would
/// otherwise have had. A device never adopts a downloaded catalog as its own,
/// so a fresh one gets its crops from the shards too.
///
/// The result must stay what every earlier build already merges: same schema,
/// same `user_version`, same page size and the same WAL flag in the header.
/// The `the_snapshot_*` tests pin those.
pub fn snapshot_for_upload(conn: &Connection, dest: &Path) -> Result<(), CatalogError> {
// The source is read through a second connection, attached to the
// snapshot's, so it needs to be a file. Every catalog is one.
let source = conn
.path()
.filter(|p| !p.is_empty())
.map(PathBuf::from)
.ok_or_else(|| CatalogError::Io("the catalog to snapshot has no file".into()))?;
// Not needed for consistency, since the read transaction below sees the
// WAL, but it keeps the live WAL from growing across syncs, as before.
checkpoint(conn)?;
// A leftover from a pass that died mid-build would otherwise be built on.
for stale in [
dest.to_path_buf(),
sidecar_of(dest, "-wal"),
sidecar_of(dest, "-journal"),
sidecar_of(dest, "-shm"),
] {
match std::fs::remove_file(&stale) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(CatalogError::Io(format!("{}: {e}", stale.display()))),
}
}
let out = Connection::open(dest)?;
build_snapshot(conn, &source, &out)?;
// This device's local album folders are paths and SAF grants nobody
// else can use; the merge never reads them, and the snapshot is what
// a fresh device would otherwise adopt whole.
out.execute_batch("DROP TABLE IF EXISTS album_folders")?;
verify_snapshot(&out)?;
Ok(())
}
/// `catalog.sqlite` + `-wal` → `catalog.sqlite-wal`.
fn sidecar_of(path: &Path, suffix: &str) -> PathBuf {
let mut s = path.as_os_str().to_owned();
s.push(suffix);
PathBuf::from(s)
}
/// Schema name the source catalog is attached under while a snapshot is built.
const SOURCE_SCHEMA: &str = "snap_src";
/// The body of [`snapshot_for_upload`]: fill the empty database `out` from
/// the catalog at `source`.
fn build_snapshot(conn: &Connection, source: &Path, out: &Connection) -> Result<(), CatalogError> {
// Settings that only take on an empty file, copied from the source so the
// result is the file a backup would have been.
let page_size: i64 = conn.query_row("PRAGMA main.page_size", [], |r| r.get(0))?;
let auto_vacuum: i64 = conn.query_row("PRAGMA main.auto_vacuum", [], |r| r.get(0))?;
out.pragma_update(None, "page_size", page_size)?;
out.pragma_update(None, "auto_vacuum", auto_vacuum)?;
// A scratch file, rebuilt whole on every pass and checked before upload:
// durability during the build buys nothing. MEMORY rather than OFF keeps
// ROLLBACK defined on the failure path.
out.pragma_update(None, "journal_mode", "MEMORY")?;
out.pragma_update(None, "synchronous", "OFF")?;
// The rows were checked when they were written, and the copy has them
// all by the end. The bundled SQLite turns foreign keys on by default,
// and with them a multi-row INSERT into `images` scans `images` for
// children of every row it adds (`shadowed_by` refers to the same table
// and has no index): 1.2 s of a 1.7 s snapshot on 24k images.
out.pragma_update(None, "foreign_keys", false)?;
// Bound as a parameter, so a path containing a quote cannot break out.
out.execute(
&format!("ATTACH DATABASE ?1 AS {SOURCE_SCHEMA}"),
[source.to_string_lossy().as_ref()],
)?;
let result = copy_schema_and_rows(out);
if let Err(e) = out.execute(&format!("DETACH DATABASE {SOURCE_SCHEMA}"), []) {
log::warn!("failed to detach the catalog from its snapshot: {e}");
}
result?;
// Last, and outside any transaction, which is the only place it can be
// set: the header says WAL, as every snapshot uploaded so far has.
out.pragma_update(None, "journal_mode", "WAL")?;
Ok(())
}
fn copy_schema_and_rows(out: &Connection) -> Result<(), CatalogError> {
let tx = out.unchecked_transaction()?;
// The first read of the source opens its read snapshot. Everything from
// here, schema included, is as of that one moment.
let objects: Vec<(String, String, String)> = {
let mut stmt = tx.prepare(&format!(
"SELECT type, name, sql FROM {SOURCE_SCHEMA}.sqlite_master
WHERE sql IS NOT NULL AND name NOT LIKE 'sqlite\\_%' ESCAPE '\\'
ORDER BY rowid"
))?;
let rows = stmt.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?;
rows.collect::<Result<_, _>>()?
};
let user_version: i64 =
tx.query_row(&format!("PRAGMA {SOURCE_SCHEMA}.user_version"), [], |r| {
r.get(0)
})?;
let application_id: i64 =
tx.query_row(&format!("PRAGMA {SOURCE_SCHEMA}.application_id"), [], |r| {
r.get(0)
})?;
// Tables and their rows first, then indexes, triggers and views, so that
// an index is built once over the data rather than maintained per row,
// and no trigger fires on the copy. Foreign keys are off on this
// connection, so the order tables are filled in does not matter.
for (_, name, sql) in objects.iter().filter(|(k, _, _)| k == "table") {
// Verbatim: an unqualified CREATE lands in `main`, the snapshot.
tx.execute_batch(sql)?;
let columns: Vec<String> = {
let mut stmt = tx.prepare("SELECT name FROM pragma_table_info(?1, 'main')")?;
let rows = stmt.query_map([name], |r| r.get::<_, String>(0))?;
rows.collect::<Result<_, _>>()?
};
let select = columns
.iter()
.map(|c| {
if name == "faces" && c == "crop" {
"NULL".to_string()
} else {
quote_ident(c)
}
})
.collect::<Vec<_>>()
.join(", ");
let insert = columns
.iter()
.map(|c| quote_ident(c))
.collect::<Vec<_>>()
.join(", ");
let table = quote_ident(name);
tx.execute(
&format!(
"INSERT INTO main.{table} ({insert})
SELECT {select} FROM {SOURCE_SCHEMA}.{table}"
),
[],
)?;
}
// AUTOINCREMENT's counters live in a table the filter above skips; the
// CREATE of such a table makes an empty one here.
let has_sequence: bool = tx.query_row(
&format!(
"SELECT EXISTS(SELECT 1 FROM {SOURCE_SCHEMA}.sqlite_master
WHERE name = 'sqlite_sequence')"
),
[],
|r| r.get(0),
)?;
if has_sequence {
tx.execute_batch(&format!(
"DELETE FROM main.sqlite_sequence;
INSERT INTO main.sqlite_sequence SELECT * FROM {SOURCE_SCHEMA}.sqlite_sequence;"
))?;
}
for (_, _, sql) in objects.iter().filter(|(k, _, _)| k != "table") {
tx.execute_batch(sql)?;
}
tx.pragma_update(None, "user_version", user_version)?;
tx.pragma_update(None, "application_id", application_id)?;
tx.commit()?;
Ok(())
}
/// `name` as an SQL identifier, whatever it contains.
fn quote_ident(name: &str) -> String {
format!("\"{}\"", name.replace('"', "\"\""))
}
/// 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.
///
/// What [`crate::recovery`] takes its NFR-R2 backups with. A backup is the
/// file the user may have to *live on*, so it keeps the face crops that
/// [`snapshot_for_upload`] leaves out, and a byte-for-byte page copy is the
/// right tool. Keeping it here beside the upload 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)
}
/// 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);
// The merge brings in rows the backfill exists for — assignments whose
// word this device has no term for, from a remote older than v6 — so the
// next open must run it, whether or not the stamp happened to move.
if let Some(path) = conn.path().filter(|p| !p.is_empty()) {
crate::backfilled::forget(Path::new(path));
}
// 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}");
}
// After every merge, because a merge is where two devices' people meet:
// the same name typed on each, or a redirect one of them made. Its own
// transaction, and a failure is logged rather than returned -- what the
// merge took is committed and valid whether or not the duplicates were
// folded, and the next pass tries again. Runs on the sync worker, never
// the UI thread, and costs ~10 ms when there is nothing to do.
if result.is_ok() {
if let Err(e) = crate::dedup_people::run(conn) {
log::warn!("dedup after the catalog merge: {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);
}
/// One image, known to the server by `file_id`, in a catalog.
fn with_image(c: &Connection, file_id: i64) -> dr_types::ImageId {
c.execute(
"INSERT OR IGNORE INTO roots(id, kind, label) VALUES (1, 'remote', 'lib')",
[],
)
.unwrap();
c.execute(
"INSERT INTO images(root_id, source_ref, added_at) VALUES (1, ?1, 0)",
[format!("IMG_{file_id}.CR3")],
)
.unwrap();
let id = c.last_insert_rowid();
c.execute(
"INSERT INTO remote(image_id, file_id) VALUES (?1, ?2)",
[id, file_id],
)
.unwrap();
dr_types::ImageId(id as u64)
}
#[test]
fn an_album_and_its_exports_reach_another_device_but_its_folder_does_not() {
use crate::albums::{self, Place};
let dir = tempdir();
let remote_path = dir.join("remote.sqlite");
let snap = dir.join("snap.sqlite");
{
// The desktop: two albums, one on the server and one on its own
// disk, each with an export of the same photograph.
let r = seeded(&dir.join("desktop.sqlite"));
let img = with_image(&r, 4242);
let web = albums::create(&r, "Web", &Place::Server("Shared/Web".into())).unwrap();
let print = albums::create(&r, "Print", &Place::Local("/mnt/print".into())).unwrap();
albums::record_exports(&r, web, &[(img, "IMG_4242.jpg".into())]).unwrap();
albums::record_exports(&r, print, &[(img, "IMG_4242.tif".into())]).unwrap();
snapshot_for_upload(&r, &snap).unwrap();
std::fs::rename(&snap, &remote_path).unwrap();
}
// The tablet knows the same file under its own image id.
let local = seeded(&dir.join("tablet.sqlite"));
with_image(&local, 1);
let img = with_image(&local, 4242);
let report = merge_remote(&local, &remote_path).unwrap();
assert_eq!(report.albums_taken, 2);
assert_eq!(report.album_exports_added, 2);
assert!(report.local_changed());
let all = albums::list(&local).unwrap();
let print = all.iter().find(|a| a.name == "Print").unwrap();
let web = all.iter().find(|a| a.name == "Web").unwrap();
assert_eq!(print.place, None, "the desktop's disk is not the tablet's");
assert_eq!(web.place, Some(Place::Server("Shared/Web".into())));
assert_eq!(albums::sources(&local, web.id).unwrap(), vec![img]);
// Nothing changed on either side, so a second pass takes nothing.
let again = merge_remote(&local, &remote_path).unwrap();
assert_eq!(again.albums_taken, 0);
assert_eq!(again.album_exports_added, 0);
}
#[test]
fn a_device_that_never_made_an_album_still_uploads_this_ones() {
use crate::albums::{self, Place};
let dir = tempdir();
let remote_path = dir.join("remote.sqlite");
{
// A snapshot from a build that predates albums altogether.
let r = seeded(&remote_path);
checkpoint(&r).unwrap();
}
let local = seeded(&dir.join("local.sqlite"));
albums::create(&local, "Web", &Place::Server("Web".into())).unwrap();
let report = merge_remote(&local, &remote_path).unwrap();
assert_eq!(report.albums_taken, 0);
// Its albums table is absent, so nothing was compared — and the local
// album has still to reach the server.
assert!(
albums::list(&local).unwrap().len() == 1,
"the local album survives a merge with a catalog that has none"
);
}
#[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");
}
/// Device-side setup for the snapshot tests: an image both devices know by
/// its cross-device file id, and one face on it carrying `crop`.
fn with_a_face(c: &Connection, face_id: i64, x: f64, crop: &[u8]) {
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 remote(image_id, file_id) VALUES (1, 5000)", [])
.unwrap();
c.execute(
"INSERT INTO faces
(id, image_id, x, y, w, h, landmarks, detector_confidence, embedding,
crop_px, model_id, detected_at, crop)
VALUES (?1, 1, ?2, 0.2, 0.2, 0.2, X'00', 0.9, X'00', 150.0, 'w600k_mbf', 0, ?3)",
rusqlite::params![face_id, x, crop],
)
.unwrap();
}
fn crop_of(c: &Connection, face_id: i64) -> Option<Vec<u8>> {
c.query_row("SELECT crop FROM faces WHERE id = ?1", [face_id], |r| {
r.get(0)
})
.unwrap()
}
/// The snapshot is built table by table rather than copied, so what has to
/// hold is that it is still the same database bar the crops: every table,
/// index and row, and the header fields an older build checks before it
/// will merge (`user_version`) or open it the way it always has (the WAL
/// flag and page size a backup-API copy carried).
#[test]
fn the_snapshot_is_the_catalog_bar_the_crops() {
let dir = tempdir();
let live = dir.join("catalog.sqlite");
let snap = dir.join("snap.sqlite");
let c = seeded(&live);
with_a_face(&c, 7, 0.3, &[7u8; 4096]);
c.execute(
"INSERT INTO collections(uuid, name, kind, created, revision, modified)
VALUES ('u1', 'Iceland', 0, 0, 1, 1)",
[],
)
.unwrap();
// Created on first use rather than by a migration: the copy must not
// depend on the migrations knowing every table.
crate::duplicates::ensure_probe_table(&c).unwrap();
snapshot_for_upload(&c, &snap).unwrap();
let out = Connection::open(&snap).unwrap();
let objects = |conn: &Connection| -> Vec<(String, String)> {
let mut stmt = conn
.prepare("SELECT type, name FROM sqlite_master ORDER BY type, name")
.unwrap();
let rows = stmt
.query_map([], |r| Ok((r.get(0).unwrap(), r.get(1).unwrap())))
.unwrap();
rows.map(Result::unwrap).collect()
};
assert_eq!(objects(&out), objects(&c), "the snapshot's schema differs");
for (kind, table) in objects(&c) {
if kind != "table" {
continue;
}
let count = |conn: &Connection| -> i64 {
conn.query_row(&format!("SELECT COUNT(*) FROM \"{table}\""), [], |r| {
r.get(0)
})
.unwrap()
};
assert_eq!(count(&out), count(&c), "rows differ in {table}");
}
for pragma in [
"user_version",
"application_id",
"page_size",
"journal_mode",
] {
let read = |conn: &Connection| -> String {
conn.query_row(&format!("PRAGMA {pragma}"), [], |r| {
r.get::<_, rusqlite::types::Value>(0)
})
.map(|v| format!("{v:?}"))
.unwrap()
};
assert_eq!(read(&out), read(&c), "{pragma} differs");
}
drop(out);
// Bytes 18 and 19 of the header are 2 for a WAL database, which is
// what every snapshot uploaded before this one said.
let header = std::fs::read(&snap).unwrap();
assert_eq!(&header[18..20], &[2, 2], "the snapshot is not WAL-flagged");
// The face is there; its pixels are not.
let out = Connection::open(&snap).unwrap();
assert_eq!(crop_of(&out, 7), None);
}
/// A pass that died mid-build leaves a file behind; the next one must
/// build afresh rather than on top of it.
#[test]
fn a_leftover_snapshot_is_replaced_not_built_on() {
let dir = tempdir();
let live = dir.join("catalog.sqlite");
let snap = dir.join("snap.sqlite");
let c = seeded(&live);
std::fs::write(&snap, b"not a database").unwrap();
snapshot_for_upload(&c, &snap).unwrap();
// And twice over a good one, which is the steady state.
snapshot_for_upload(&c, &snap).unwrap();
let out = Connection::open(&snap).unwrap();
let v: i64 = out
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(v, schema::SCHEMA_VERSION);
}
/// The merge side of a crop-less snapshot: the other device's names still
/// cross over — the match is by box, not by pixels — and this device's
/// own crop is left exactly as it was, never replaced by the snapshot's
/// NULL.
#[test]
fn merging_a_crop_less_snapshot_keeps_local_crops_and_takes_the_names() {
let dir = tempdir();
let snap = dir.join("snap.sqlite");
let desktop = seeded(&dir.join("desktop.sqlite"));
with_a_face(&desktop, 42, 0.31, &[1u8; 3000]);
desktop
.execute(
"INSERT INTO people(id, uuid, name, ignored, created, revision, modified)
VALUES (3, 'u-anna', 'Anna', 0, 0, 1, 1)",
[],
)
.unwrap();
desktop
.execute(
"INSERT INTO face_person(face_id, person_id, probability, confirmed)
VALUES (42, 3, 0.9, 1)",
[],
)
.unwrap();
snapshot_for_upload(&desktop, &snap).unwrap();
// The tablet found the same face itself, under its own row id, and has
// its own crop of it — from its own detection or from the shards.
let tablet = seeded(&dir.join("tablet.sqlite"));
let mine = vec![9u8; 2500];
with_a_face(&tablet, 7, 0.30, &mine);
let report = merge_remote(&tablet, &snap).unwrap();
assert_eq!(report.people_inserted, 1);
assert_eq!(report.faces_assigned, 1);
let named: (String, bool) = tablet
.query_row(
"SELECT p.name, fp.confirmed
FROM face_person fp JOIN people p ON p.id = fp.person_id
WHERE fp.face_id = 7",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(named, ("Anna".to_string(), true));
assert_eq!(
crop_of(&tablet, 7),
Some(mine),
"the merge touched a local crop"
);
// Idempotent over a crop-less remote too.
let again = merge_remote(&tablet, &snap).unwrap();
assert!(!again.local_changed());
}
}