diff --git a/core/dr-catalog/src/face_shard.rs b/core/dr-catalog/src/face_shard.rs index 65d50e8..82ef910 100644 --- a/core/dr-catalog/src/face_shard.rs +++ b/core/dr-catalog/src/face_shard.rs @@ -30,6 +30,7 @@ //! identity every client agrees on (FR-NC-5), and it survives a server-side //! move, so a shard written before a reorganisation still applies after it. +use std::collections::{HashMap, HashSet}; use std::path::{Path, PathBuf}; use rusqlite::{Connection, OptionalExtension}; @@ -155,6 +156,45 @@ impl FaceShardStore { .flatten() } + /// [`indexed_at`](Self::indexed_at) for every entry at once. + /// + /// What a pass over the whole library asks instead of one lookup per image: + /// the export compares every marker in the catalog with this, and a lookup + /// each was 19,000 statements prepared and run on every sync pass that had + /// nothing to send. + fn all_indexed_at(&self) -> Result>, CatalogError> { + let mut q = self + .index + .prepare("SELECT file_id, model_id, indexed_at FROM entries")?; + let rows = q.query_map([], |r| { + Ok(( + (r.get::<_, i64>(0)? as u64, r.get::<_, String>(1)?), + r.get::<_, Option>(2)?, + )) + })?; + Ok(rows.collect::>()?) + } + + /// [`held_model`](Self::held_model) for every file at once: one statement, + /// ordered exactly as that one is, keeping the first row per file. + fn all_held_models(&self, model_id: &str) -> Result, CatalogError> { + let mut q = self.index.prepare(&format!( + "SELECT file_id, model_id FROM entries + WHERE {} = ?1 + ORDER BY file_id, indexed_at DESC NULLS LAST, model_id", + crate::faces::embedder_sql("model_id") + ))?; + let rows = q.query_map([crate::faces::embedder_of(model_id)], |r| { + Ok((r.get::<_, i64>(0)? as u64, r.get::<_, String>(1)?)) + })?; + let mut out = HashMap::new(); + for row in rows { + let (file, model) = row?; + out.entry(file).or_insert(model); + } + Ok(out) + } + /// The pipeline this store holds an image under, among those sharing /// `model_id`'s embedder — the most recently indexed where a peer has /// sent more than one. @@ -784,6 +824,18 @@ pub fn export_to_shards_reporting( /// enough that the reporting is lost in the write it accompanies. const REPORT_EVERY: usize = 25; + // What the store holds, read once. A put below rewrites only its own + // file's entries -- its generation, and siblings it supersedes -- so a + // file already written in this pass is asked of the store again and every + // other answer is the one a lookup would have given. + // + // An index that cannot be read answers as each lookup did: nothing held. + let held = store.all_indexed_at().unwrap_or_else(|e| { + log::debug!("reading the shard index: {e}"); + HashMap::new() + }); + let mut written: HashSet = HashSet::new(); + let total = rows.len(); let mut exported = 0; for (seen, (file_id, image_id, edge, indexed_at, model_id)) in rows.into_iter().enumerate() { @@ -800,10 +852,14 @@ pub fn export_to_shards_reporting( // // The comparison is against when the *catalog* indexed it, so a // re-index is visible and an unchanged image still costs nothing. - if store - .indexed_at(file_id as u64, model_id) - .is_some_and(|was| was >= indexed_at) - { + let was = if written.contains(&(file_id as u64)) { + store.indexed_at(file_id as u64, model_id) + } else { + held.get(&(file_id as u64, model_id.to_string())) + .copied() + .flatten() + }; + if was.is_some_and(|was| was >= indexed_at) { continue; } let mut fq = conn.prepare( @@ -840,6 +896,7 @@ pub fn export_to_shards_reporting( &faces, Some(indexed_at), )?; + written.insert(file_id as u64); exported += 1; } progress(total, total); @@ -902,6 +959,18 @@ pub fn import_from_shards( /// the lock, waits a fraction of a second and not the whole import. const CHUNK: usize = 100; + // Every file's held pipeline, read once rather than asked per candidate -- + // 23,000 prepared lookups on every pass, nearly all of them for images + // this device already holds. Nothing below changes which pipeline the + // store holds a file under (`set_indexed_at` touches only a file already + // decided), and each candidate is a different file, so these are the + // answers the lookups gave. + // An index that cannot be read answers as each lookup did: nothing held. + let held_models = store.all_held_models(model_id).unwrap_or_else(|e| { + log::debug!("reading the shard index: {e}"); + HashMap::new() + }); + let mut adopted = 0; let mut tx = conn.unchecked_transaction()?; let mut in_chunk = 0; @@ -911,7 +980,7 @@ pub fn import_from_shards( tx = conn.unchecked_transaction()?; in_chunk = 0; } - let Some(held) = store.held_model(file_id as u64, model_id) else { + let Some(held) = held_models.get(&(file_id as u64)).cloned() else { continue; }; if let Some(local) = local {