From 0729dfa359e3d0c164a9e0794d57e0d4ec150fe9 Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Fri, 28 Aug 2026 21:19:56 +0200 Subject: [PATCH] Stop the face sync opening a database per photograph MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit "Checking faces…" never finished. `export_to_shards` asks, for every indexed image in the library, whether the shard store already holds that image at that index time — and `indexed_at` answered by opening the shard database, running its six-statement schema batch and two `pragma_table_info` queries, then querying. Once per image. 9,849 times for this library, on every sync pass, before a single face had been written. The index time now lives in the store's `index.sqlite` alongside the shard number, so the question is one indexed lookup on a connection that is already open. It stays in the shard as well — that copy is the one that travels — but nothing reads it from there on the hot path. The write side had the same shape: `put_image` opened the shard afresh for each image, which mattered little when exports were a handful of new photographs and matters a great deal now that a re-index sends thousands. The handle is kept and reused, invalidated by shard id so sealing a full one and moving to the next drops it without anything having to remember to. `INDEX_SCHEMA` is `CREATE ... IF NOT EXISTS` like the shard schema, so the new column is added on open for an index already on disk — the same trap, caught the same way. Two tests: that the index time survives reopening the store, and that an index written before the column can still be opened and written to. Co-Authored-By: Claude Opus 5 (1M context) --- core/dr-catalog/src/face_shard.rs | 128 +++++++++++++++++++++++++----- 1 file changed, 108 insertions(+), 20 deletions(-) diff --git a/core/dr-catalog/src/face_shard.rs b/core/dr-catalog/src/face_shard.rs index b127e79..7bfdde0 100644 --- a/core/dr-catalog/src/face_shard.rs +++ b/core/dr-catalog/src/face_shard.rs @@ -90,6 +90,13 @@ pub struct FaceShardStore { dir: PathBuf, index: Connection, client: String, + /// The shard currently being written to, held open across a run of writes. + /// + /// A bulk export writes thousands of images in a row, nearly all into the + /// same shard, and opening it per image means a file open and a six + /// statement schema batch each time. Held here, that cost is paid once per + /// shard instead of once per photograph. + writer: Option<(u32, Connection)>, } impl FaceShardStore { @@ -97,11 +104,18 @@ impl FaceShardStore { std::fs::create_dir_all(dir).map_err(|e| CatalogError::Io(e.to_string()))?; let index = Connection::open(dir.join("index.sqlite"))?; index.execute_batch(INDEX_SCHEMA)?; + // `INDEX_SCHEMA` is `CREATE ... IF NOT EXISTS` like the shard schema, + // so a column added to it never reaches an index already on disk. The + // same trap, caught the same way -- see `upgrade_shard`. + if !has_column(&index, "entries", "indexed_at")? { + index.execute_batch("ALTER TABLE entries ADD COLUMN indexed_at INTEGER")?; + } let client = mint_client_id(&index)?; Ok(Self { dir: dir.to_path_buf(), index, client, + writer: None, }) } @@ -119,26 +133,16 @@ impl FaceShardStore { /// existed. Both mean "cannot vouch for it", and the export treats them the /// same: send it again. pub fn indexed_at(&self, file_id: u64, model_id: &str) -> Option { - let shard: i64 = self - .index + self.index .query_row( - "SELECT shard FROM entries WHERE file_id = ?1 AND model_id = ?2", + "SELECT indexed_at FROM entries WHERE file_id = ?1 AND model_id = ?2", rusqlite::params![file_id as i64, model_id], - |r| r.get(0), + |r| r.get::<_, Option>(0), ) .optional() .ok() - .flatten()?; - let conn = self.open_shard(shard as u32, false).ok()?; - conn.query_row( - "SELECT indexed_at FROM indexed WHERE file_id = ?1 AND model_id = ?2", - rusqlite::params![file_id as i64, model_id], - |r| r.get::<_, Option>(0), - ) - .optional() - .ok() - .flatten() - .flatten() + .flatten() + .flatten() } pub fn contains(&self, file_id: u64, model_id: &str) -> bool { @@ -202,7 +206,11 @@ impl FaceShardStore { .sum::() + 64; let shard = self.active_shard(incoming)?; - let conn = self.open_shard(shard, true)?; + self.ensure_writer(shard)?; + // Both this and `self.index` below are immutable borrows of `self`, so + // they coexist; the cache is filled first, above, where the mutable + // borrow is over by the time this one starts. + let conn = &self.writer.as_ref().expect("just opened").1; let tx = conn.unchecked_transaction()?; // Replace rather than append: re-indexing an image must not double its @@ -250,11 +258,19 @@ impl FaceShardStore { tx.commit()?; self.index.execute( - "INSERT INTO entries (file_id, model_id, shard, bytes) - VALUES (?1, ?2, ?3, ?4) + "INSERT INTO entries (file_id, model_id, shard, bytes, indexed_at) + VALUES (?1, ?2, ?3, ?4, ?5) ON CONFLICT(file_id, model_id) DO UPDATE SET - shard = excluded.shard, bytes = excluded.bytes", - rusqlite::params![file_id as i64, model_id, shard as i64, incoming as i64], + shard = excluded.shard, + bytes = excluded.bytes, + indexed_at = excluded.indexed_at", + rusqlite::params![ + file_id as i64, + model_id, + shard as i64, + incoming as i64, + indexed_at + ], )?; self.index.execute( "INSERT INTO faces_meta (file_id, model_id) VALUES (?1, ?2) @@ -268,6 +284,21 @@ impl FaceShardStore { Ok(shard) } + /// Open the shard being written to, reusing the handle where we already + /// have it. + /// + /// Invalidated by id rather than by a flag, so sealing a full shard and + /// moving to the next one drops the old handle without anything having to + /// remember to. + fn ensure_writer(&mut self, shard: u32) -> Result<(), CatalogError> { + if self.writer.as_ref().map(|(id, _)| *id) == Some(shard) { + return Ok(()); + } + let conn = self.open_shard(shard, true)?; + self.writer = Some((shard, conn)); + Ok(()) + } + /// The shard currently open for writing, sealing and starting a new one /// where this batch would overflow the cap. fn active_shard(&self, incoming: u64) -> Result { @@ -691,6 +722,15 @@ CREATE TABLE IF NOT EXISTS entries ( model_id TEXT NOT NULL, shard INTEGER NOT NULL, bytes INTEGER NOT NULL, + -- When the catalog indexed this image, so `export_to_shards` can tell a + -- re-indexed image from an unchanged one. + -- + -- Here as well as in the shard, and this is the copy that gets read. The + -- export asks about every image in the library on every pass, and reading + -- it from the shard meant opening a shard database per image — 9,849 file + -- opens and schema batches for this library, which is not a slow sync but + -- a hung one. This index is already open. + indexed_at INTEGER, PRIMARY KEY(file_id, model_id) ); CREATE INDEX IF NOT EXISTS entries_shard ON entries(shard); @@ -1317,4 +1357,52 @@ mod catalog_round_trip { assert!(faces[0].crop.is_empty(), "a crop was invented from nowhere"); let _ = std::fs::remove_dir_all(&dir); } + + /// The index time is read on every image of every sync, so it lives in the + /// index rather than in the shard. Reading it from the shard meant opening + /// a shard database per image — 9,849 file opens and schema batches for the + /// reference library, which is not a slow sync but a hung one. + #[test] + fn the_index_time_survives_reopening_the_store() { + let dir = tempdir("index-time"); + { + let mut s = FaceShardStore::open(&dir).unwrap(); + s.put_image_at(4, "w600k_mbf", 1024, &[old_face(4, 1)], Some(9_000)) + .unwrap(); + assert_eq!(s.indexed_at(4, "w600k_mbf"), Some(9_000)); + } + // Reopened from disk, which is what every sync after the first does. + let s = FaceShardStore::open(&dir).unwrap(); + assert_eq!(s.indexed_at(4, "w600k_mbf"), Some(9_000)); + assert_eq!(s.indexed_at(5, "w600k_mbf"), None, "an image never stored"); + let _ = std::fs::remove_dir_all(&dir); + } + + /// An index written before the column existed is the ordinary case on any + /// machine that has ever synced — the same `CREATE IF NOT EXISTS` trap the + /// shards fell into, and it has to be caught here too. + #[test] + fn an_index_from_before_the_column_is_upgraded_on_open() { + let dir = tempdir("old-index"); + { + // The pre-column shape, written by hand. + let c = Connection::open(dir.join("index.sqlite")).unwrap(); + c.execute_batch( + "CREATE TABLE entries ( + file_id INTEGER NOT NULL, + model_id TEXT NOT NULL, + shard INTEGER NOT NULL, + bytes INTEGER NOT NULL, + PRIMARY KEY(file_id, model_id) + );", + ) + .unwrap(); + } + + let mut s = FaceShardStore::open(&dir).expect("opening an old index failed"); + s.put_image_at(6, "w600k_mbf", 1024, &[old_face(6, 2)], Some(11)) + .expect("writing to an upgraded index failed"); + assert_eq!(s.indexed_at(6, "w600k_mbf"), Some(11)); + let _ = std::fs::remove_dir_all(&dir); + } }