Merge: origin's sync ordering and adoption work, its scan-complete trigger ported into library_ui/

This commit is contained in:
2026-09-21 11:32:05 +02:00
13 changed files with 432 additions and 43 deletions
+178 -9
View File
@@ -178,6 +178,67 @@ impl FaceShardStore {
.flatten()
}
/// The other pipelines this file is held under that share `model_id`'s
/// embedder — the generations a put of `model_id` may supersede.
fn siblings(&self, file_id: u64, model_id: &str) -> Vec<String> {
let mut stmt = match self.index.prepare(&format!(
"SELECT model_id FROM entries
WHERE file_id = ?1 AND model_id != ?2 AND {} = ?3",
crate::faces::embedder_sql("model_id")
)) {
Ok(s) => s,
Err(_) => return Vec::new(),
};
stmt.query_map(
rusqlite::params![
file_id as i64,
model_id,
crate::faces::embedder_of(model_id)
],
|r| r.get::<_, String>(0),
)
.map(|rows| rows.filter_map(|r| r.ok()).collect())
.unwrap_or_default()
}
/// Whether a pass this file is already held under outranks `model_id`,
/// so a put of `model_id` would add a generation nobody would adopt.
pub fn outranked(&self, file_id: u64, model_id: &str) -> bool {
use dr_types::FaceDetector;
let Some(incoming) = FaceDetector::for_model_id(model_id) else {
return false;
};
self.siblings(file_id, model_id)
.iter()
.filter_map(|m| FaceDetector::for_model_id(m))
.any(|held| held.outranks(incoming))
}
/// Forget the index entries for generations of this file that `model_id`
/// outranks. The bytes stay where they are — a sealed shard is
/// immutable — but the store stops offering them, and a later export or
/// merge writes nothing for them again.
fn supersede(&self, file_id: u64, model_id: &str) -> Result<(), CatalogError> {
use dr_types::FaceDetector;
let Some(incoming) = FaceDetector::for_model_id(model_id) else {
return Ok(());
};
for held in self.siblings(file_id, model_id) {
let weaker = FaceDetector::for_model_id(&held).is_some_and(|h| incoming.outranks(h));
if weaker {
self.index.execute(
"DELETE FROM entries WHERE file_id = ?1 AND model_id = ?2",
rusqlite::params![file_id as i64, held],
)?;
self.index.execute(
"DELETE FROM faces_meta WHERE file_id = ?1 AND model_id = ?2",
rusqlite::params![file_id as i64, held],
)?;
}
}
Ok(())
}
pub fn contains(&self, file_id: u64, model_id: &str) -> bool {
self.index
.query_row(
@@ -233,6 +294,15 @@ impl FaceShardStore {
faces: &[SharedFace],
indexed_at: Option<i64>,
) -> Result<u32, CatalogError> {
// One generation per image per embedder. A store carried every pass
// — 24,123 entries for 19,089 images on the reference library, a
// third of its 293 MB — and only the strongest was ever adopted.
// A weaker pass arriving after a stronger one is not written; a
// stronger one arriving retires the weaker from the index.
if self.outranked(file_id, model_id) {
return Ok(0);
}
self.supersede(file_id, model_id)?;
let incoming = faces
.iter()
.map(|f| BYTES_PER_FACE + if f.crop.is_empty() { 0 } else { BYTES_PER_CROP })
@@ -474,15 +544,28 @@ impl FaceShardStore {
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
)?;
let mut q =
src.prepare("SELECT file_id, model_id, faces_found, source_edge FROM indexed")?;
let images: Vec<(i64, String, i64, i64)> = q
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))?
// The peer's marker travels with the image: it is what lets
// `import_from_shards` record the adoption under the time the peer
// indexed it, and so what keeps `export_to_shards` from reading the
// adoption as a re-index and sending the peer's faces back out under
// this device's name. A shard from before the column has none.
let mut q = src.prepare(&format!(
"SELECT file_id, model_id, faces_found, source_edge, {} FROM indexed",
match has_column(&src, "indexed", "indexed_at") {
Ok(true) => "indexed_at",
_ => "NULL",
}
))?;
let images: Vec<(i64, String, i64, i64, Option<i64>)> = q
.query_map([], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?))
})?
.collect::<Result<_, _>>()?;
let mut adopted = 0;
for (file_id, model_id, _found, edge) in images {
if self.contains(file_id as u64, &model_id) {
for (file_id, model_id, _found, edge, indexed_at) in images {
if self.contains(file_id as u64, &model_id) || self.outranked(file_id as u64, &model_id)
{
continue;
}
let mut fq = src.prepare(&format!(
@@ -501,12 +584,27 @@ impl FaceShardStore {
let faces: Vec<SharedFace> = fq
.query_map(rusqlite::params![file_id, &model_id], read_shared_face)?
.collect::<Result<_, _>>()?;
self.put_image(file_id as u64, &model_id, edge as u32, &faces)?;
self.put_image_at(file_id as u64, &model_id, edge as u32, &faces, indexed_at)?;
adopted += 1;
}
Ok(adopted)
}
/// Record when the catalog indexed a held image, for an entry that
/// arrived without a marker — a peer's shard from before the column.
pub fn set_indexed_at(
&self,
file_id: u64,
model_id: &str,
at: i64,
) -> Result<(), CatalogError> {
self.index.execute(
"UPDATE entries SET indexed_at = ?3 WHERE file_id = ?1 AND model_id = ?2",
rusqlite::params![file_id as i64, model_id, at],
)?;
Ok(())
}
/// Read back everything held for one image.
pub fn get_image(
&self,
@@ -798,8 +896,21 @@ pub fn import_from_shards(
})?
.collect::<Result<_, _>>()?;
/// Images per write transaction. Large enough that fourteen thousand
/// adoptions are a hundred and forty commits rather than fourteen
/// thousand; small enough that a read on the UI thread, queued behind
/// the lock, waits a fraction of a second and not the whole import.
const CHUNK: usize = 100;
let mut adopted = 0;
let mut tx = conn.unchecked_transaction()?;
let mut in_chunk = 0;
for (file_id, image_id, local) in candidates {
if in_chunk == CHUNK {
tx.commit()?;
tx = conn.unchecked_transaction()?;
in_chunk = 0;
}
let Some(held) = store.held_model(file_id as u64, model_id) else {
continue;
};
@@ -858,15 +969,42 @@ pub fn import_from_shards(
})
.collect();
crate::faces::record_detections(
conn,
crate::faces::record_detections_within(
&tx,
dr_types::ImageId(image_id as u64),
&held,
edge,
&local,
)?;
// The peer's marker, not this moment. `record_detections` stamps the
// run as now, and `export_to_shards` reads a marker newer than the
// shard's as a re-index — so every adopted image went straight back
// out as this device's own work: 14,100 adopted, 15,457 "newly
// indexed" on the next pass, and twenty-two shards of a peer's faces
// uploaded again under a second name. Where the peer's shard carried
// no marker, the store takes the catalog's, so the two agree either
// way and the export sees nothing to send.
match store.indexed_at(file_id as u64, &held) {
Some(theirs) => {
tx.execute(
"UPDATE face_index SET indexed_at = ?3
WHERE image_id = ?1 AND model_id = ?2",
rusqlite::params![image_id, held, theirs],
)?;
}
None => {
let ours: i64 = tx.query_row(
"SELECT indexed_at FROM face_index WHERE image_id = ?1 AND model_id = ?2",
rusqlite::params![image_id, held],
|r| r.get(0),
)?;
store.set_indexed_at(file_id as u64, &held, ours)?;
}
}
adopted += 1;
in_chunk += 1;
}
tx.commit()?;
Ok(adopted)
}
@@ -1102,6 +1240,32 @@ mod tests {
assert!(!s.contains(1, "lvface"));
}
/// One generation per image per embedder: a stronger detector's pass
/// retires a weaker one from the index, and a weaker pass arriving after
/// a stronger is not written at all.
#[test]
fn a_stronger_pass_retires_a_weaker_one_and_a_weaker_is_not_added() {
let dir = tempdir();
let mut s = FaceShardStore::open(&dir).unwrap();
s.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap();
s.put_image(1, "scrfd_10g+w600k_mbf", 1024, &[face(1, 2)])
.unwrap();
assert!(s.contains(1, "scrfd_10g+w600k_mbf"));
assert!(!s.contains(1, "w600k_mbf"), "the fast pass was not retired");
assert_eq!(s.len(), 1, "faces_meta still counts the retired pass");
s.put_image(1, "scrfd_2.5g+w600k_mbf", 1024, &[face(1, 3)])
.unwrap();
assert!(
!s.contains(1, "scrfd_2.5g+w600k_mbf"),
"a weaker pass was added"
);
assert_eq!(
s.held_model(1, "w600k_mbf").as_deref(),
Some("scrfd_10g+w600k_mbf")
);
}
#[test]
fn re_storing_an_image_replaces_rather_than_doubling_it() {
let dir = tempdir();
@@ -1322,6 +1486,11 @@ mod catalog_round_trip {
assert!((got[0].landmarks[2].0 - 0.15).abs() < 1e-5);
let emb = faces::embeddings(&b, "w600k_mbf").unwrap();
assert!(emb.iter().any(|e| e.embedding[0] == 1));
// And what B adopted is not B's work: its next export sends nothing.
// Adopting used to stamp the run as now, so every adopted image went
// back out under B's name as a re-index.
assert_eq!(export_to_shards(&b, &mut store_b, "w600k_mbf").unwrap(), 0);
}
/// The desktop switched to a stronger detector part-way through the
+19 -2
View File
@@ -279,7 +279,25 @@ pub fn record_detections(
faces: &[DetectedFace],
) -> Result<Vec<FaceId>, CatalogError> {
let tx = conn.unchecked_transaction()?;
let ids = record_detections_within(&tx, image_id, model_id, source_edge, faces)?;
tx.commit()?;
Ok(ids)
}
/// [`record_detections`] inside a transaction the caller owns.
///
/// For a caller recording many images at once — the shard import adopts
/// fourteen thousand in one pass — where a commit per image is fourteen
/// thousand fsyncs and fourteen thousand turns at the write lock that every
/// read on the UI thread queues behind. `unchecked_transaction` cannot nest,
/// so the batching has to be offered here rather than wrapped from above.
pub fn record_detections_within(
tx: &Connection,
image_id: ImageId,
model_id: &str,
source_edge: u32,
faces: &[DetectedFace],
) -> Result<Vec<FaceId>, CatalogError> {
// Everything the old faces knew, so it can be carried across the
// replacement. Read only when there is something to carry it onto: a
// pass that found nothing has nothing to match, and decoding a vector
@@ -287,7 +305,7 @@ pub fn record_detections(
let prior = if faces.is_empty() {
Vec::new()
} else {
read_priors(&tx, image_id)?
read_priors(tx, image_id)?
};
tx.execute("DELETE FROM faces WHERE image_id = ?1", [image_id.0 as i64])?;
@@ -389,7 +407,6 @@ pub fn record_detections(
],
)?;
tx.commit()?;
Ok(ids)
}
+111
View File
@@ -119,6 +119,8 @@ pub struct MergeReport {
pub keywords_fused: usize,
/// Keyword assignments taken from the remote.
pub keywords_assigned: usize,
/// Images whose capture metadata was taken from the remote.
pub metadata_adopted: usize,
}
impl MergeReport {
@@ -133,6 +135,7 @@ impl MergeReport {
|| self.keywords_deleted > 0
|| self.keywords_fused > 0
|| self.keywords_assigned > 0
|| self.metadata_adopted > 0
}
/// Whether the local catalog holds anything the remote did not, and so
@@ -193,10 +196,67 @@ pub fn merge_all(conn: &Connection) -> Result<MergeReport, CatalogError> {
merge_collections_within(&tx, &mut report)?;
merge_keywords_within(&tx, &mut report)?;
merge_people_within(&tx, &mut report)?;
merge_metadata_within(&tx, &mut report)?;
tx.commit()?;
Ok(report)
}
/// Adopt capture metadata from an attached catalog, on its own.
pub fn merge_metadata(conn: &Connection) -> Result<MergeReport, CatalogError> {
let tx = conn.unchecked_transaction()?;
let mut report = MergeReport::default();
merge_metadata_within(&tx, &mut report)?;
tx.commit()?;
Ok(report)
}
/// Capture metadata a peer's sweep already read, for images this device has
/// not dated yet.
///
/// The `images` table is local state and the merge leaves it alone — except
/// for these columns, which are not: a capture time, an offset, a camera, a
/// lens and an ISO are facts about the file's bytes, identical on every
/// device, and read by fetching a header per image across the whole library
/// (`dr_ui::library::spawn_sweep`). A fresh device inherits its peers'
/// thumbnails and faces from the shards and then spent hours re-reading
/// every header for the timeline; the snapshot it had just merged held
/// every one of those dates.
///
/// Matched by `oc:fileid`, as collection membership is. Only rows still at
/// `metadata_state < 2` take anything, and only from a remote row at 2: a
/// date this device read for itself is never overwritten, and a peer that
/// has not read one has nothing to give. The sweep's own query
/// (`metadata_state < 2`) then finds nothing left to do for them.
const METADATA_BY_FILE_ID: &str = "
UPDATE main.images
SET captured_at = r.captured_at,
captured_offset = coalesce(main.images.captured_offset, r.captured_offset),
camera = coalesce(main.images.camera, r.camera),
lens = coalesce(main.images.lens, r.lens),
iso = coalesce(main.images.iso, r.iso),
metadata_state = 2
FROM (SELECT lr.image_id, ri.captured_at, ri.captured_offset,
ri.camera, ri.lens, ri.iso
FROM remote_cat.images ri
JOIN remote_cat.remote rr ON rr.image_id = ri.id
JOIN main.remote lr ON lr.file_id = rr.file_id
WHERE ri.metadata_state >= 2 AND ri.captured_at IS NOT NULL) AS r
WHERE main.images.id = r.image_id
AND main.images.metadata_state < 2";
fn merge_metadata_within(tx: &Connection, report: &mut MergeReport) -> Result<(), CatalogError> {
// A snapshot from before these columns, or from a library with no server
// behind it, has nothing to join on.
if !remote_has(tx, "remote")?
|| !remote_has_column(tx, "images", "metadata_state")?
|| !remote_has_column(tx, "images", "captured_offset")?
{
return Ok(());
}
report.metadata_adopted = tx.execute(METADATA_BY_FILE_ID, [])?;
Ok(())
}
/// Merge people and identity judgements from an attached catalog.
///
/// The people half of [`merge_all`], on its own, for the same reason the other
@@ -1168,6 +1228,57 @@ mod tests {
// ---- integration over two real catalogs ------------------------------
/// A fresh device takes the capture dates a peer's sweep read, matched by
/// `oc:fileid`, and never overwrites a date it read for itself.
#[test]
fn capture_metadata_arrives_for_undated_images_only() {
let c = two_catalogs();
// Three photographs on both devices: 1 undated here and dated there;
// 2 dated on both, differently; 3 undated on both.
for id in 1..=3 {
add_image_without_hash(&c, "main", id);
add_image_without_hash(&c, "remote_cat", id + 10);
add_remote_id(&c, "main", id, 100 + id);
add_remote_id(&c, "remote_cat", id + 10, 100 + id);
}
c.execute(
"UPDATE remote_cat.images
SET captured_at = 1000, captured_offset = 60, camera = 'X', metadata_state = 2
WHERE id = 11",
[],
)
.unwrap();
c.execute(
"UPDATE remote_cat.images SET captured_at = 2000, metadata_state = 2 WHERE id = 12",
[],
)
.unwrap();
c.execute(
"UPDATE main.images SET captured_at = 2222, metadata_state = 2 WHERE id = 2",
[],
)
.unwrap();
let report = merge_metadata(&c).unwrap();
assert_eq!(report.metadata_adopted, 1);
let row = |id: i64| -> (Option<i64>, Option<i64>, Option<String>, i64) {
c.query_row(
"SELECT captured_at, captured_offset, camera, metadata_state
FROM main.images WHERE id = ?1",
[id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
)
.unwrap()
};
assert_eq!(row(1), (Some(1000), Some(60), Some("X".into()), 2));
assert_eq!(row(2), (Some(2222), None, None, 2));
assert_eq!(row(3), (None, None, None, 0));
// Idempotent: a second pass finds nothing left to take.
assert_eq!(merge_metadata(&c).unwrap().metadata_adopted, 0);
}
fn two_catalogs() -> Connection {
attached_remote(schema::for_attached("remote_cat"))
}