Merge: a face sync that can actually complete

Three faults between a laptop that had indexed a library and a tablet that only
ever received part of it.

The face store was the one part of the catalog still on the rollback journal at
synchronous=FULL — 21.3 ms a commit against the 0.05 ms everything else pays,
four commits per photograph, on the order of fourteen minutes of pure fsync for
ten thousand images. It now writes the way the catalog and the thumbnail store
do, with a checkpoint before upload so the file that goes to the server is
complete without its write-ahead log.

Every idle pass read 94 MB of shards to decide it had nothing to send; the name
and the size answer that.

And the 60-second total request timeout was a floor on link speed rather than a
hang detector: a 25 MB shard needed a sustained 425 KB/s or it failed, and then
retried and failed again indefinitely. It now times out on inactivity.

The merge itself was never at fault — a probe that stands up an empty catalog,
adopts the shards and merges the real one in gets all 611 of a person's faces,
not 200.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-08-28 23:08:29 +02:00
co-authored by Claude Opus 5
4 changed files with 148 additions and 23 deletions
+103 -3
View File
@@ -103,6 +103,7 @@ impl FaceShardStore {
pub fn open(dir: &Path) -> Result<Self, CatalogError> {
std::fs::create_dir_all(dir).map_err(|e| CatalogError::Io(e.to_string()))?;
let index = Connection::open(dir.join("index.sqlite"))?;
fast_writes(&index)?;
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
@@ -257,7 +258,12 @@ impl FaceShardStore {
)?;
tx.commit()?;
self.index.execute(
// One transaction, not three. Each of these was autocommitting, and an
// autocommit is a durable write — so a single photograph cost four
// commits counting the shard's own, and an export of ten thousand paid
// forty thousand of them.
let ix = self.index.unchecked_transaction()?;
ix.execute(
"INSERT INTO entries (file_id, model_id, shard, bytes, indexed_at)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT(file_id, model_id) DO UPDATE SET
@@ -272,15 +278,16 @@ impl FaceShardStore {
indexed_at
],
)?;
self.index.execute(
ix.execute(
"INSERT INTO faces_meta (file_id, model_id) VALUES (?1, ?2)
ON CONFLICT(file_id, model_id) DO NOTHING",
rusqlite::params![file_id as i64, model_id],
)?;
self.index.execute(
ix.execute(
"UPDATE shards SET bytes = bytes + ?2 WHERE id = ?1",
rusqlite::params![shard as i64, incoming as i64],
)?;
ix.commit()?;
Ok(shard)
}
@@ -341,11 +348,32 @@ impl FaceShardStore {
return Err(CatalogError::Io(format!("face shard {shard} is missing")));
}
let conn = Connection::open(&path)?;
fast_writes(&conn)?;
conn.execute_batch(SHARD_SCHEMA)?;
upgrade_shard(&conn)?;
Ok(conn)
}
/// Fold every write-ahead log back into the database files.
///
/// **A shard is uploaded by reading its file.** Under WAL the most recent
/// commits live in a `-wal` sidecar until a checkpoint moves them, so
/// reading the `.sqlite` alone would ship a database missing exactly the
/// faces just written — and a peer would adopt it and see nothing wrong.
/// The sync calls this before it reads anything.
///
/// `TRUNCATE` rather than the default passive checkpoint: passive gives up
/// when a reader holds the log, which would leave the same gap while
/// reporting success.
pub fn checkpoint(&self) -> Result<(), CatalogError> {
self.index
.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)")?;
if let Some((_, conn)) = self.writer.as_ref() {
conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)")?;
}
Ok(())
}
pub fn shard_path(&self, shard: u32) -> PathBuf {
self.dir
.join(format!("shard-{}-{shard:04}.sqlite", self.client))
@@ -469,6 +497,27 @@ impl FaceShardStore {
}
}
/// Put a connection into the mode the rest of the catalog already uses.
///
/// The face store was the one place still on SQLite's default rollback journal
/// at `synchronous = FULL`, where the catalog (`schema::configure`) and the
/// thumbnail store both run WAL at `NORMAL`. Measured on this project's own
/// filesystem the difference is **21.3 ms per commit against 0.05 ms** — four
/// hundred times — and an export commits several times per photograph, so ten
/// thousand images spent something like fourteen minutes doing nothing but
/// waiting for fsync. That was the "checking faces…" that never finished.
///
/// `NORMAL` rather than `FULL` is the same trade the rest of the catalog makes:
/// a shard is derived data, and the cost of losing the last commit to a power
/// cut is that the next pass exports that image again.
fn fast_writes(conn: &Connection) -> Result<(), CatalogError> {
// A `journal_mode` change returns the new mode as a row, so it has to be
// queried rather than executed.
conn.query_row("PRAGMA journal_mode = WAL", [], |_| Ok(()))?;
conn.execute_batch("PRAGMA synchronous = NORMAL")?;
Ok(())
}
/// Add the columns a shard written by an older build is missing.
///
/// [`SHARD_SCHEMA`] is entirely `CREATE ... IF NOT EXISTS`, which does exactly
@@ -1431,4 +1480,55 @@ mod catalog_round_trip {
assert_eq!(s.indexed_at(6, "w600k_mbf"), Some(11));
let _ = std::fs::remove_dir_all(&dir);
}
/// The face store was the one part of the catalog still on the rollback
/// journal at `synchronous = FULL`, which cost 21 ms a commit where the
/// rest pays 0.05 ms. An export commits several times per photograph.
#[test]
fn the_store_writes_the_way_the_rest_of_the_catalog_does() {
let dir = tempdir("wal");
let mut s = FaceShardStore::open(&dir).unwrap();
s.put_image(1, "w600k_mbf", 1024, &[old_face(1, 1)])
.unwrap();
for db in [dir.join("index.sqlite"), s.shard_path(0)] {
let c = Connection::open(&db).unwrap();
let mode: String = c
.query_row("PRAGMA journal_mode", [], |r| r.get(0))
.unwrap();
assert_eq!(mode, "wal", "{} is not in WAL", db.display());
}
let _ = std::fs::remove_dir_all(&dir);
}
/// **The property the upload depends on.** A shard is sent by reading its
/// file; under WAL the newest commits sit in a `-wal` sidecar that is not
/// sent with it. Without a checkpoint the server would receive a database
/// missing exactly the faces just exported, and a peer would adopt it and
/// see nothing wrong.
#[test]
fn a_checkpointed_shard_stands_alone_without_its_write_ahead_log() {
let dir = tempdir("checkpoint");
let mut s = FaceShardStore::open(&dir).unwrap();
for id in 1..=5u64 {
s.put_image(id, "w600k_mbf", 1024, &[old_face(id, id as u8)])
.unwrap();
}
s.checkpoint().unwrap();
// Copy *only* the database, exactly as the upload reads it.
let sent = dir.join("as-uploaded.sqlite");
std::fs::copy(s.shard_path(0), &sent).unwrap();
let c = Connection::open_with_flags(
&sent,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
.unwrap();
let faces: i64 = c
.query_row("SELECT COUNT(*) FROM faces", [], |r| r.get(0))
.unwrap();
assert_eq!(faces, 5, "the uploaded file is missing recent writes");
let _ = std::fs::remove_dir_all(&dir);
}
}
+16 -4
View File
@@ -637,11 +637,23 @@ pub fn http_client(user_agent: &str) -> Result<reqwest::Client, RemoteError> {
// webpki-roots set is still installed alongside.
.tls_certs_only(extra_roots())
// A request that hangs forever is indistinguishable from a worker that
// died, and cost a long time to tell apart once. Connect and total
// timeouts turn that into an error the UI can show. Generous enough for
// a slow phone on mobile data; the login poll has its own deadline.
// died, and cost a long time to tell apart once. These turn that into
// an error the UI can show.
.connect_timeout(std::time::Duration::from_secs(15))
.timeout(std::time::Duration::from_secs(60))
// **Inactivity, not duration.** This was a 60-second *total* timeout,
// which is not a hang detector at all — it is a floor on link speed.
// A face shard runs to 25 MB, so it demanded a sustained 425 KB/s or
// the transfer failed; and having failed it was retried on the next
// pass, and failed again, for ever. A tablet on ordinary wifi could
// therefore never finish adopting a library's faces, and nothing said
// why: each attempt looked like a network blip rather than an
// arithmetic impossibility.
//
// `read_timeout` fires when *no bytes arrive* for the given period,
// which is the condition actually worth failing on. A slow transfer
// that is still moving now finishes, however long it takes, while a
// connection that has genuinely died is still caught in a minute.
.read_timeout(std::time::Duration::from_secs(60))
.build()
.map_err(|e| RemoteError::Network(e.to_string()))
}
+1 -1
View File
@@ -87,7 +87,7 @@ _None._
| FR-EXP-9 | [`core/dr-decode/src/lib.rs:506`](../core/dr-decode/src/lib.rs#L506), [`core/dr-export/src/lib.rs:128`](../core/dr-export/src/lib.rs#L128), [`core/dr-export/src/lib.rs:1`](../core/dr-export/src/lib.rs#L1), [`core/dr-gpu/src/adjust.rs:1068`](../core/dr-gpu/src/adjust.rs#L1068), [`ui/dr-ui/src/develop.rs:2829`](../ui/dr-ui/src/develop.rs#L2829), [`ui/dr-ui/src/lib.rs:362`](../ui/dr-ui/src/lib.rs#L362) |
| FR-NC-1 | [`core/dr-sync-nextcloud/src/auth.rs:132`](../core/dr-sync-nextcloud/src/auth.rs#L132), [`core/dr-sync-nextcloud/src/auth.rs:44`](../core/dr-sync-nextcloud/src/auth.rs#L44), [`core/dr-sync-nextcloud/src/session.rs:128`](../core/dr-sync-nextcloud/src/session.rs#L128), [`ui/dr-ui/src/launch.rs:256`](../ui/dr-ui/src/launch.rs#L256), [`ui/dr-ui/src/launch.rs:49`](../ui/dr-ui/src/launch.rs#L49), [`ui/dr-ui/src/launch_ui.rs:344`](../ui/dr-ui/src/launch_ui.rs#L344) |
| FR-NC-10 | [`ui/dr-ui/src/export.rs:1`](../ui/dr-ui/src/export.rs#L1), [`ui/dr-ui/src/lib.rs:424`](../ui/dr-ui/src/lib.rs#L424), [`ui/dr-ui/src/library.rs:1049`](../ui/dr-ui/src/library.rs#L1049), [`ui/dr-ui/src/library.rs:1693`](../ui/dr-ui/src/library.rs#L1693), [`ui/dr-ui/src/library.rs:544`](../ui/dr-ui/src/library.rs#L544), [`ui/dr-ui/src/library.rs:846`](../ui/dr-ui/src/library.rs#L846), [`ui/dr-ui/src/library_ui.rs:1607`](../ui/dr-ui/src/library_ui.rs#L1607), [`ui/dr-ui/src/library_ui.rs:3365`](../ui/dr-ui/src/library_ui.rs#L3365), [`ui/dr-ui/src/library_ui.rs:499`](../ui/dr-ui/src/library_ui.rs#L499), [`ui/dr-ui/src/sidecar_cache.rs:1`](../ui/dr-ui/src/sidecar_cache.rs#L1) |
| FR-NC-12 | [`core/dr-sync-nextcloud/src/lib.rs:34`](../core/dr-sync-nextcloud/src/lib.rs#L34), [`core/dr-sync-nextcloud/src/lib.rs:892`](../core/dr-sync-nextcloud/src/lib.rs#L892), [`core/dr-sync/src/lib.rs:157`](../core/dr-sync/src/lib.rs#L157), [`core/dr-sync/src/lib.rs:40`](../core/dr-sync/src/lib.rs#L40), [`core/dr-sync/src/reachability.rs:1`](../core/dr-sync/src/reachability.rs#L1), [`ui/dr-ui/src/remote.rs:1`](../ui/dr-ui/src/remote.rs#L1) |
| FR-NC-12 | [`core/dr-sync-nextcloud/src/lib.rs:34`](../core/dr-sync-nextcloud/src/lib.rs#L34), [`core/dr-sync-nextcloud/src/lib.rs:904`](../core/dr-sync-nextcloud/src/lib.rs#L904), [`core/dr-sync/src/lib.rs:157`](../core/dr-sync/src/lib.rs#L157), [`core/dr-sync/src/lib.rs:40`](../core/dr-sync/src/lib.rs#L40), [`core/dr-sync/src/reachability.rs:1`](../core/dr-sync/src/reachability.rs#L1), [`ui/dr-ui/src/remote.rs:1`](../ui/dr-ui/src/remote.rs#L1) |
| FR-NC-2 | [`core/dr-sync-nextcloud/src/session.rs:128`](../core/dr-sync-nextcloud/src/session.rs#L128), [`core/dr-sync-nextcloud/src/session.rs:34`](../core/dr-sync-nextcloud/src/session.rs#L34), [`platform/dr-plat/src/secrets.rs:82`](../platform/dr-plat/src/secrets.rs#L82) |
| FR-NC-3 | [`core/dr-decode/src/locate.rs:1`](../core/dr-decode/src/locate.rs#L1), [`core/dr-decode/src/preview.rs:148`](../core/dr-decode/src/preview.rs#L148), [`core/dr-sync/src/capability.rs:41`](../core/dr-sync/src/capability.rs#L41), [`core/dr-thumbs/src/lib.rs:1`](../core/dr-thumbs/src/lib.rs#L1), [`ui/dr-ui/src/library.rs:1`](../ui/dr-ui/src/library.rs#L1), [`ui/dr-ui/src/library.rs:2724`](../ui/dr-ui/src/library.rs#L2724), [`ui/dr-ui/src/library.rs:3195`](../ui/dr-ui/src/library.rs#L3195), [`ui/dr-ui/src/library_ui.rs:1`](../ui/dr-ui/src/library_ui.rs#L1), [`ui/dr-ui/src/library_ui.rs:3627`](../ui/dr-ui/src/library_ui.rs#L3627), [`ui/dr-ui/src/library_ui.rs:4593`](../ui/dr-ui/src/library_ui.rs#L4593), [`ui/dr-ui/ui/app.slint:316`](../ui/dr-ui/ui/app.slint#L316), [`ui/dr-ui/ui/settings.slint:372`](../ui/dr-ui/ui/settings.slint#L372), [`ui/dr-ui/ui/settings.slint:72`](../ui/dr-ui/ui/settings.slint#L72) |
| FR-NC-4 | [`core/dr-sync-nextcloud/src/propfind.rs:100`](../core/dr-sync-nextcloud/src/propfind.rs#L100), [`core/dr-sync-nextcloud/src/propfind.rs:51`](../core/dr-sync-nextcloud/src/propfind.rs#L51), [`core/dr-sync/src/capability.rs:6`](../core/dr-sync/src/capability.rs#L6), [`core/dr-sync/src/lib.rs:157`](../core/dr-sync/src/lib.rs#L157), [`core/dr-sync/src/scan.rs:93`](../core/dr-sync/src/scan.rs#L93), [`ui/dr-ui/src/launch.rs:49`](../ui/dr-ui/src/launch.rs#L49) |
+28 -15
View File
@@ -391,16 +391,41 @@ async fn sync_face_shards(
.unwrap_or_default();
// ---- upload ----------------------------------------------------------
//
// Every commit since the last pass is still in a write-ahead log, and a
// shard is uploaded by reading its file — so without this the upload would
// ship a database missing precisely the faces just exported.
if let Err(e) = store.checkpoint() {
log::warn!("face sync: checkpointing the shard store: {e}");
}
let local = store.shards().map_err(|e| e.to_string())?;
for (n, shard) in local.iter().enumerate() {
let name = shard_name(&client, shard.id);
let path = store.shard_path(shard.id);
// **Decided before the file is read.** A sealed shard the server
// already has is byte-identical by construction, and the client is in
// the name so nobody else could have written it — the name alone
// settles it. Reading first meant every idle sync pulled ninety-four
// megabytes off disk to conclude it had nothing to send.
let on_disk = std::fs::metadata(&path).map(|m| m.len()).unwrap_or(0);
let skip = match remote.get(&name) {
Some(_) if shard.sealed => true,
Some(size) => *size == on_disk,
None => false,
};
if skip {
continue;
}
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.
// single one is a visible wait on any connection. Announced after the
// skip, or an idle pass claims to be sending five shards and sends
// none; and before the put, because the wait is the upload.
let _ = tx.send(SyncMessage::Status(format!(
"sending faces: shard {}/{} ({} MB)",
n + 1,
@@ -408,18 +433,6 @@ async fn sync_face_shards(
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
// again whenever its size differs, which is the only way it changes.
let skip = match remote.get(&name) {
Some(_) if shard.sealed => true,
Some(size) => *size == bytes.len() as u64,
None => false,
};
if skip {
continue;
}
let target = RemotePath::new(format!("{}/{name}", face_base.as_str()));
match backend.put(&target, bytes, None).await {
Ok(_) => report.face_shards_uploaded += 1,