Merge: a face sync that finishes
Build and test / Desktop (Linux) (push) Successful in 35m30s
Build and test / Layer separation (push) Successful in 53s
Traceability / Requirement traces (push) Successful in 38s
🐳 Android image / Build and push (push) Successful in 1s
Build and test / android-image (push) Successful in 2s
Build and test / Android (aarch64) (push) Failing after 58m37s
Build and test / Desktop (Linux) (push) Successful in 35m30s
Build and test / Layer separation (push) Successful in 53s
Traceability / Requirement traces (push) Successful in 38s
🐳 Android image / Build and push (push) Successful in 1s
Build and test / android-image (push) Successful in 2s
Build and test / Android (aarch64) (push) Failing after 58m37s
The export asked whether each image was already in the shard store at its current index time, and answered by opening the shard database — schema batch, pragma probes and all — once per photograph. 9,849 opens per sync pass for this library, before any face was written, which showed up as a sync stuck on "checking faces…" and never coming back. The index time moves to the store's own index, which is already open, and the write handle is held across a run instead of reopened per image. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -90,6 +90,13 @@ pub struct FaceShardStore {
|
|||||||
dir: PathBuf,
|
dir: PathBuf,
|
||||||
index: Connection,
|
index: Connection,
|
||||||
client: String,
|
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 {
|
impl FaceShardStore {
|
||||||
@@ -97,11 +104,18 @@ impl FaceShardStore {
|
|||||||
std::fs::create_dir_all(dir).map_err(|e| CatalogError::Io(e.to_string()))?;
|
std::fs::create_dir_all(dir).map_err(|e| CatalogError::Io(e.to_string()))?;
|
||||||
let index = Connection::open(dir.join("index.sqlite"))?;
|
let index = Connection::open(dir.join("index.sqlite"))?;
|
||||||
index.execute_batch(INDEX_SCHEMA)?;
|
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)?;
|
let client = mint_client_id(&index)?;
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
dir: dir.to_path_buf(),
|
dir: dir.to_path_buf(),
|
||||||
index,
|
index,
|
||||||
client,
|
client,
|
||||||
|
writer: None,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -119,26 +133,16 @@ impl FaceShardStore {
|
|||||||
/// existed. Both mean "cannot vouch for it", and the export treats them the
|
/// existed. Both mean "cannot vouch for it", and the export treats them the
|
||||||
/// same: send it again.
|
/// same: send it again.
|
||||||
pub fn indexed_at(&self, file_id: u64, model_id: &str) -> Option<i64> {
|
pub fn indexed_at(&self, file_id: u64, model_id: &str) -> Option<i64> {
|
||||||
let shard: i64 = self
|
self.index
|
||||||
.index
|
|
||||||
.query_row(
|
.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],
|
rusqlite::params![file_id as i64, model_id],
|
||||||
|r| r.get(0),
|
|r| r.get::<_, Option<i64>>(0),
|
||||||
)
|
)
|
||||||
.optional()
|
.optional()
|
||||||
.ok()
|
.ok()
|
||||||
.flatten()?;
|
.flatten()
|
||||||
let conn = self.open_shard(shard as u32, false).ok()?;
|
.flatten()
|
||||||
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<i64>>(0),
|
|
||||||
)
|
|
||||||
.optional()
|
|
||||||
.ok()
|
|
||||||
.flatten()
|
|
||||||
.flatten()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn contains(&self, file_id: u64, model_id: &str) -> bool {
|
pub fn contains(&self, file_id: u64, model_id: &str) -> bool {
|
||||||
@@ -202,7 +206,11 @@ impl FaceShardStore {
|
|||||||
.sum::<u64>()
|
.sum::<u64>()
|
||||||
+ 64;
|
+ 64;
|
||||||
let shard = self.active_shard(incoming)?;
|
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()?;
|
let tx = conn.unchecked_transaction()?;
|
||||||
// Replace rather than append: re-indexing an image must not double its
|
// Replace rather than append: re-indexing an image must not double its
|
||||||
@@ -250,11 +258,19 @@ impl FaceShardStore {
|
|||||||
tx.commit()?;
|
tx.commit()?;
|
||||||
|
|
||||||
self.index.execute(
|
self.index.execute(
|
||||||
"INSERT INTO entries (file_id, model_id, shard, bytes)
|
"INSERT INTO entries (file_id, model_id, shard, bytes, indexed_at)
|
||||||
VALUES (?1, ?2, ?3, ?4)
|
VALUES (?1, ?2, ?3, ?4, ?5)
|
||||||
ON CONFLICT(file_id, model_id) DO UPDATE SET
|
ON CONFLICT(file_id, model_id) DO UPDATE SET
|
||||||
shard = excluded.shard, bytes = excluded.bytes",
|
shard = excluded.shard,
|
||||||
rusqlite::params![file_id as i64, model_id, shard as i64, incoming as i64],
|
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(
|
self.index.execute(
|
||||||
"INSERT INTO faces_meta (file_id, model_id) VALUES (?1, ?2)
|
"INSERT INTO faces_meta (file_id, model_id) VALUES (?1, ?2)
|
||||||
@@ -268,6 +284,21 @@ impl FaceShardStore {
|
|||||||
Ok(shard)
|
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
|
/// The shard currently open for writing, sealing and starting a new one
|
||||||
/// where this batch would overflow the cap.
|
/// where this batch would overflow the cap.
|
||||||
fn active_shard(&self, incoming: u64) -> Result<u32, CatalogError> {
|
fn active_shard(&self, incoming: u64) -> Result<u32, CatalogError> {
|
||||||
@@ -691,6 +722,15 @@ CREATE TABLE IF NOT EXISTS entries (
|
|||||||
model_id TEXT NOT NULL,
|
model_id TEXT NOT NULL,
|
||||||
shard INTEGER NOT NULL,
|
shard INTEGER NOT NULL,
|
||||||
bytes 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)
|
PRIMARY KEY(file_id, model_id)
|
||||||
);
|
);
|
||||||
CREATE INDEX IF NOT EXISTS entries_shard ON entries(shard);
|
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");
|
assert!(faces[0].crop.is_empty(), "a crop was invented from nowhere");
|
||||||
let _ = std::fs::remove_dir_all(&dir);
|
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);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user