Merge: a face sync that says what it is doing
Build and test / Desktop (Linux) (push) Successful in 39m22s
Build and test / Layer separation (push) Successful in 52s
🐳 Android image / Build and push (push) Successful in 4s
Build and test / android-image (push) Successful in 5s
Traceability / Requirement traces (push) Successful in 53s
Build and test / Android (aarch64) (push) Failing after 53m19s
Build and test / Desktop (Linux) (push) Successful in 39m22s
Build and test / Layer separation (push) Successful in 52s
🐳 Android image / Build and push (push) Successful in 4s
Build and test / android-image (push) Successful in 5s
Traceability / Requirement traces (push) Successful in 53s
Build and test / Android (aarch64) (push) Failing after 53m19s
The face pass announced itself once and then went quiet for the whole export, upload and adopt. On a first export after a re-index that is eight minutes of a motionless progress bar, which reads as a hang. It now reports the images it is preparing, the shards it is sending and their size, and the shards it is taking in. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -531,6 +531,23 @@ pub fn export_to_shards(
|
||||
conn: &Connection,
|
||||
store: &mut FaceShardStore,
|
||||
model_id: &str,
|
||||
) -> Result<usize, CatalogError> {
|
||||
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<usize, CatalogError> {
|
||||
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::<Result<_, _>>()?;
|
||||
|
||||
/// 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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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<SyncMessage>,
|
||||
) -> 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 {
|
||||
|
||||
Reference in New Issue
Block a user