diff --git a/core/dr-catalog/src/sync.rs b/core/dr-catalog/src/sync.rs index 6b5efea..010538d 100644 --- a/core/dr-catalog/src/sync.rs +++ b/core/dr-catalog/src/sync.rs @@ -46,18 +46,206 @@ pub fn checkpoint(conn: &Connection) -> Result<(), CatalogError> { Ok(()) } -/// Write a consistent snapshot of the catalog to `dest`, ready to upload. +/// Write a consistent snapshot of the catalog to `dest`, ready to upload, +/// without the face crops. /// -/// 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. +/// 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> { - let out = copy_to(conn, dest)?; - strip_face_crops(&out)?; + // 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)?; 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::>()? + }; + 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 = { + 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::>()? + }; + let select = columns + .iter() + .map(|c| { + if name == "faces" && c == "crop" { + "NULL".to_string() + } else { + quote_ident(c) + } + }) + .collect::>() + .join(", "); + let insert = columns + .iter() + .map(|c| quote_ident(c)) + .collect::>() + .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`. /// @@ -82,12 +270,12 @@ fn verify_snapshot(snapshot: &Connection) -> Result<(), CatalogError> { /// 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. +/// 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 { checkpoint(conn)?; @@ -103,39 +291,6 @@ pub(crate) fn copy_to(conn: &Connection, dest: &Path) -> Result 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 @@ -380,4 +535,188 @@ mod tests { .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> { + 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()); + } } diff --git a/docs/dev/catalog.md b/docs/dev/catalog.md index b2fbaa5..ae442d4 100644 --- a/docs/dev/catalog.md +++ b/docs/dev/catalog.md @@ -698,9 +698,23 @@ pair that had no other home. **A WAL database is not one file.** Committed transactions can sit in `catalog.sqlite-wal` with the main file lagging, so copying `catalog.sqlite` alone uploads a torn snapshot — internally consistent -as of some older point, silently missing everything since. Upload therefore runs a `TRUNCATE` -checkpoint and then SQLite's backup API, which serialises against concurrent writers rather than -racing them. It never copies the live file. +as of some older point, silently missing everything since. The upload therefore never copies the +live file. It builds the snapshot in an empty file. It attaches the catalog, creates each table +from the catalog's own schema, and fills it with `INSERT … SELECT`, all inside one transaction. +That transaction holds a single read snapshot of the catalog, so concurrent writers are serialised +rather than raced, as the backup API did before. The file is then checked with `quick_check` before +it goes anywhere. + +**The face crops stay out of the upload.** A crop is a ~5 KB JPEG on each `faces` row. On a 19k-face +library they are 96 MB of a 158 MB catalog. The face shards carry them to other devices, once each. +The merge reads a remote face's box and model to match it to a local one, never its pixels. No +device adopts a downloaded catalog as its own: a fresh device starts empty and takes faces, crops +included, from the shards. So the snapshot's `crop` is NULL, and a merge never writes a local +crop. They were first stripped (2026-08) by copying the whole file with the backup API, setting +`crop` to NULL and `VACUUM`ing. That wrote the file about three times to upload 50 MB. Since #71 +the snapshot is built without them. Its header is kept as it was (`user_version`, page size and +the WAL flag), so every earlier build merges it unchanged. NFR-R2 backups still use the backup API +and keep the crops, because a backup is a file the user may have to live on. **Integer primary keys are not identities.** Two devices each allocate `collections.id = 1` for different collections, so a row-level merge keyed on the integer id would collide them. Collections