diff --git a/core/dr-catalog/src/face_shard.rs b/core/dr-catalog/src/face_shard.rs index eea6a15..bf2ddf6 100644 --- a/core/dr-catalog/src/face_shard.rs +++ b/core/dr-catalog/src/face_shard.rs @@ -103,6 +103,7 @@ impl FaceShardStore { pub fn open(dir: &Path) -> Result { 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); + } } diff --git a/core/dr-sync-nextcloud/src/lib.rs b/core/dr-sync-nextcloud/src/lib.rs index 8db855e..7cda5e9 100644 --- a/core/dr-sync-nextcloud/src/lib.rs +++ b/core/dr-sync-nextcloud/src/lib.rs @@ -637,11 +637,23 @@ pub fn http_client(user_agent: &str) -> Result { // 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())) } diff --git a/docs/traceability.md b/docs/traceability.md index e4b5cac..e1c2e53 100644 --- a/docs/traceability.md +++ b/docs/traceability.md @@ -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) | diff --git a/ui/dr-ui/src/derived_sync.rs b/ui/dr-ui/src/derived_sync.rs index dec6393..7eeff5a 100644 --- a/ui/dr-ui/src/derived_sync.rs +++ b/ui/dr-ui/src/derived_sync.rs @@ -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,