diff --git a/core/dr-catalog/src/face_shard.rs b/core/dr-catalog/src/face_shard.rs index 7bfdde0..eea6a15 100644 --- a/core/dr-catalog/src/face_shard.rs +++ b/core/dr-catalog/src/face_shard.rs @@ -531,6 +531,23 @@ pub fn export_to_shards( conn: &Connection, store: &mut FaceShardStore, model_id: &str, +) -> Result { + export_to_shards_reporting(conn, store, model_id, &mut |_, _| {}) +} + +/// [`export_to_shards`], reporting how far it has got. +/// +/// `progress` is called with `(done, total)` as images are written. The first +/// export after a re-index is thousands of photographs and the better part of +/// ten minutes, and a sync that says nothing for that long is indistinguishable +/// from one that has hung — which is exactly how it was reported. Called every +/// few dozen images rather than every one, so the channel behind it is not the +/// expensive part of the loop. +pub fn export_to_shards_reporting( + conn: &Connection, + store: &mut FaceShardStore, + model_id: &str, + progress: &mut dyn FnMut(usize, usize), ) -> Result { let mut q = conn.prepare( "SELECT r.file_id, fi.image_id, fi.source_edge, fi.indexed_at @@ -545,8 +562,16 @@ pub fn export_to_shards( })? .collect::>()?; + /// How often to report. Small enough that a bar moves visibly, large + /// enough that the reporting is lost in the write it accompanies. + const REPORT_EVERY: usize = 25; + + let total = rows.len(); let mut exported = 0; - for (file_id, image_id, edge, indexed_at) in rows { + for (seen, (file_id, image_id, edge, indexed_at)) in rows.into_iter().enumerate() { + if seen.is_multiple_of(REPORT_EVERY) { + progress(seen, total); + } // **Not `contains`.** Asking only whether the shard has heard of this // image means it is exported once and never again — so re-indexing it // updates the catalog and nothing else, and every other device keeps @@ -593,6 +618,7 @@ pub fn export_to_shards( )?; exported += 1; } + progress(total, total); Ok(exported) } diff --git a/ui/dr-ui/src/derived_sync.rs b/ui/dr-ui/src/derived_sync.rs index f9eac64..dec6393 100644 --- a/ui/dr-ui/src/derived_sync.rs +++ b/ui/dr-ui/src/derived_sync.rs @@ -158,7 +158,7 @@ async fn run( sync_shards(backend, &base, thumbs_dir, scratch, &mut report).await?; let _ = tx.send(SyncMessage::Status("checking faces…".into())); - sync_face_shards(backend, &base, catalog_path, scratch, &mut report).await?; + sync_face_shards(backend, &base, catalog_path, scratch, &mut report, tx).await?; let _ = tx.send(SyncMessage::Status("checking collections…".into())); sync_catalog(backend, &base, catalog_path, scratch, &mut report).await?; @@ -328,6 +328,7 @@ async fn sync_face_shards( catalog_path: &Path, scratch: &Path, report: &mut SyncReport, + tx: &std::sync::mpsc::Sender, ) -> Result<(), String> { use dr_catalog::face_shard::{self, FaceShardStore}; @@ -351,7 +352,21 @@ async fn sync_face_shards( // ---- everything this device has detected, into the shards ------------ if let Ok(catalog) = dr_catalog::Catalog::open(catalog_path) { - match face_shard::export_to_shards(catalog.connection(), &mut store, model) { + // The first export after a re-index walks the whole library and writes + // thousands of images. Saying so is the difference between a sync that + // looks slow and one that looks broken — this pass sat on "checking + // faces…" for eight minutes and was reported as a hang. + let mut say = |done: usize, total: usize| { + let _ = tx.send(SyncMessage::Status(format!( + "preparing faces to send: {done}/{total}" + ))); + }; + match face_shard::export_to_shards_reporting( + catalog.connection(), + &mut store, + model, + &mut say, + ) { Ok(0) => {} Ok(n) => log::info!("face sync: {n} newly indexed image(s) ready to upload"), Err(e) => log::warn!("face sync: exporting to shards: {e}"), @@ -377,12 +392,21 @@ async fn sync_face_shards( // ---- upload ---------------------------------------------------------- let local = store.shards().map_err(|e| e.to_string())?; - for shard in &local { + for (n, shard) in local.iter().enumerate() { let path = store.shard_path(shard.id); let Ok(bytes) = std::fs::read(&path) else { continue; }; let name = shard_name(&client, shard.id); + // Face shards carry crops and run to tens of megabytes each, so a + // single one is a visible wait on any connection. Announced before the + // put rather than after, because the wait is the upload. + let _ = tx.send(SyncMessage::Status(format!( + "sending faces: shard {}/{} ({} MB)", + n + 1, + local.len(), + bytes.len() / 1_048_576 + ))); // Sealed and present means byte-identical, and the client is in the // name so nobody else could have written it. The open shard goes up @@ -411,6 +435,12 @@ async fn sync_face_shards( if owner == client || store.has_adopted(name, *size) { continue; } + // Taking in a peer's shard means a download and then a row-by-row + // merge, both of which take a while on a shard carrying crops. + let _ = tx.send(SyncMessage::Status(format!( + "taking in faces from another device ({} MB)", + size / 1_048_576 + ))); let source = RemotePath::new(format!("{}/{name}", face_base.as_str())); let bytes = match backend.get(&RemoteId::Path(source), None).await {