From 12065f9adc54eeb90668e40958ce717f6971657e Mon Sep 17 00:00:00 2001 From: flupkede Date: Wed, 16 Sep 2026 23:13:24 +0200 Subject: [PATCH] [worker] fix: auto-recover LMDB format corruption via sequential wipe + rebuild Symbol-rebuild failures of the MDB_BAD_VALSIZE class (data written by an older storage-major, observed on all C# repos after the arroy 0.8/heed 0.22 upgrades) now queue the repo for automatic recovery: evict stores (remove_repo's sequence, minus unregister), wipe the DB dir with the bounded lock-retry, and force-reindex through the TUI machinery whose store-open path recreates fresh formats. Recoveries are processed strictly one repo at a time (single worker, flag+queue under one lock so no lost wake-up); read-only repos are skipped with a pointer to the owning writer. Detection helper is unit-tested; queue dedup + worker start pinned by test. --- CHANGELOG.md | 4 + src/serve/mod.rs | 215 +++++++++++++++++++++++++++++++++++++++++++-- src/serve/tests.rs | 47 ++++++++++ src/serve/tui.rs | 4 +- 4 files changed, 263 insertions(+), 7 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 63848388..5fa88708 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -44,6 +44,10 @@ finalized in place with a date — no renaming/migration step needed. - **Small majors batch: dirs 7, sha2 0.11, scip 0.10, sysinfo 0.39; dead `tower`/`tower-http` direct deps removed.** dirs/scip/sysinfo were drop-in. sha2 0.11's digest arrays no longer implement `LowerHex`, so the two hash-to-hex sites (`file_meta.rs`, `chunker/mod.rs`) hex-encode the digest bytes explicitly — output unchanged. `tower` and `tower-http` were declared as direct dependencies but never imported anywhere (CORS/trace middleware never wired in); removing them shrinks the direct dependency surface (both remain in the lock transitively via axum/reqwest/hf-hub, which is upstream's business). +### Fixed + +- **Serve auto-recovers LMDB storage-format corruption with a sequential wipe + rebuild.** After the arroy 0.5→0.8 / heed 0.20→0.22 major upgrades, every repo whose on-disk database was written by the previous binary failed its symbol rebuild with `MDB_BAD_VALSIZE: Unsupported size of key/DB name/data, or wrong DUPFIXED size` (observed on all C# repos after deploy). When a symbol rebuild now fails with that error class, serve wipes the repo's DB directory — closing the LMDB envs first via the same eviction sequence `remove_repo` uses, with the same bounded retry for transient Windows lock holders — and force-reindexes it through the existing force-reindex machinery, whose store-open path recreates everything on the new formats. Recoveries are queued and processed strictly **one repo at a time**: each rebuild runs a full CPU-bound embed pass, so parallel recoveries would thrash the machine. Read-only repos are skipped with a pointer to the owning writer. This closes the gap the tantivy FTS graceful reset (above) already covered on the FTS side — the vector/symbol stores now self-heal the same upgrade boundary instead of staying red until an operator force-reindexes by hand. + ## [1.3.19] ### Changed diff --git a/src/serve/mod.rs b/src/serve/mod.rs index c1b43f6b..8bca5806 100644 --- a/src/serve/mod.rs +++ b/src/serve/mod.rs @@ -230,6 +230,18 @@ pub(crate) struct ServeState { /// drop first. The token is stored alongside the handle so `remove_repo` /// can cancel regardless of the repo's `RepoState` variant. index_tasks: DashMap, CancellationToken)>, + /// Aliases detected with LMDB storage-format corruption (e.g. + /// `MDB_BAD_VALSIZE` after a storage-layer major upgrade such as + /// arroy 0.5→0.8 / heed 0.20→0.22), queued for a wipe + full force + /// reindex. Processed strictly one at a time by the recovery worker — + /// each rebuild runs a full CPU-bound embed pass, so parallel + /// recoveries would thrash the machine. See + /// [`Self::enqueue_format_recovery`] / [`Self::recover_repo_format`]. + format_recovery_queue: std::sync::Mutex>, + /// Guarantees at most one recovery worker is alive. A worker that finds + /// the queue empty flips this back to `false` under the queue lock, so an + /// enqueue racing the worker's exit re-spawns cleanly (no lost wake-up). + format_recovery_worker_started: std::sync::atomic::AtomicBool, /// Loaded repos config (alias → path). config: std::sync::RwLock, /// Last observed mtime of the repos config file. @@ -364,6 +376,8 @@ impl ServeState { open_locks: DashMap::new(), fsw_tasks: DashMap::new(), index_tasks: DashMap::new(), + format_recovery_queue: std::sync::Mutex::new(std::collections::VecDeque::new()), + format_recovery_worker_started: std::sync::atomic::AtomicBool::new(false), config: std::sync::RwLock::new(config), config_mtime: std::sync::RwLock::new(None), config_path_override, @@ -613,6 +627,182 @@ impl ServeState { self.active_reindexes.remove(alias); } + /// True iff an error chain indicates LMDB storage-format corruption — + /// data written by an older storage-major (arroy/heed) that the current + /// one refuses to read — rather than a transient or unrelated failure. + fn is_lmdb_format_corruption(msg: &str) -> bool { + let m = msg.to_ascii_lowercase(); + m.contains("mdb_bad_valsize") + || m.contains("unsupported size of key") + || m.contains("wrong dupfixed size") + } + + /// Queue `alias` for a sequential wipe + force reindex after LMDB format + /// corruption was detected. Deduplicates; spawns the single recovery + /// worker on the first enqueue. + fn enqueue_format_recovery(self: &Arc, alias: &str) { + { + let mut queue = self + .format_recovery_queue + .lock() + .expect("format_recovery_queue lock poisoned"); + if queue.iter().any(|a| a == alias) { + return; + } + queue.push_back(alias.to_string()); + } + // Swap AFTER the push so the worker-exit path (which flips the flag + // back to `false` while still holding the queue lock) can never race + // us into a lost wake-up: either we observe `true` and the live + // worker picks up the fresh entry, or we flip `false→true` and spawn. + if !self + .format_recovery_worker_started + .swap(true, std::sync::atomic::Ordering::AcqRel) + { + let state = Arc::clone(self); + tokio::spawn(state.format_recovery_worker()); + } + } + + /// Pops queued aliases and recovers them ONE AT A TIME until the queue + /// runs dry, then exits (a later enqueue restarts a worker). + async fn format_recovery_worker(self: Arc) { + loop { + let alias = { + let mut queue = self + .format_recovery_queue + .lock() + .expect("format_recovery_queue lock poisoned"); + match queue.pop_front() { + Some(a) => a, + None => { + // Flip the flag while still holding the queue lock so + // a concurrent enqueue cannot interleave between the + // empty pop and the flag reset (lost wake-up). + self.format_recovery_worker_started + .store(false, std::sync::atomic::Ordering::Release); + return; + } + } + }; + if let Err(e) = self.recover_repo_format(&alias).await { + tracing::error!("🔧 Format recovery failed for '{}': {}", alias, e); + } + } + } + + /// Wipe + force reindex one repo whose on-disk storage was written by an + /// older storage-major. Mirrors `remove_repo`'s eviction sequence (stop + /// FSW → evict → await watcher/index shutdowns) but keeps the alias + /// registered; then deletes the DB directory (bounded retry for transient + /// Windows lock holders) and reuses the TUI force-reindex machinery — its + /// `try_open_stores` path recreates fresh stores when the directory is + /// gone, so the rebuild lands on the new arroy/heed formats. + async fn recover_repo_format(self: &Arc, alias: &str) -> Result<(), String> { + let project_path = { + let config = self + .config + .read() + .map_err(|_| "config lock poisoned".to_string())?; + if config.repo_read_only.get(alias) == Some(&true) { + return Err(format!( + "'{}' is marked read-only; rebuild its index on the owning writer", + alias + )); + } + config + .resolve(alias) + .ok_or_else(|| format!("unknown alias '{}'", alias))? + }; + let db_path = project_path.join(DB_DIR_NAME); + + // Evict in-memory holders so the LMDB env closes before the delete + // (Windows refuses to delete mmap'd files). Same order as remove_repo. + { + let _stores = self.stop_fsw(alias); + } + self.repos.remove(alias); + self.last_access.remove(alias); + self.await_fsw_shutdown(alias).await; + self.await_index_task(alias).await; + + let deadline = + Instant::now() + Duration::from_secs(crate::constants::DB_DELETE_RETRY_BUDGET_SECS); + let mut backoff_ms = crate::constants::DB_DELETE_RETRY_INITIAL_MS; + loop { + match std::fs::remove_dir_all(&db_path) { + Ok(()) => break, + Err(e) if e.kind() == std::io::ErrorKind::NotFound || !db_path.exists() => break, + Err(e) if Self::is_db_locked_error(&e) && Instant::now() < deadline => { + tracing::debug!( + "Format recovery: DB dir for '{}' still locked, retrying: {}", + alias, + e + ); + tokio::time::sleep(Duration::from_millis(backoff_ms)).await; + backoff_ms = (backoff_ms * 2).min(2_000); + } + Err(e) => { + return Err(format!( + "could not wipe {} after corruption: {}", + db_path.display(), + e + )) + } + } + } + tracing::info!( + "🔧 Format recovery: wiped stale-format DB dir for '{}' — rebuilding", + alias + ); + + match tui::spawn_force_reindex(alias.to_string(), self) { + tui::ReindexLaunch::Started => { + // Sequential guarantee: wait until this alias stops indexing + // before the worker loop picks the next one. Poll — the + // active-reindexes entry can go stale (MAX_INDEXING_SECS) on + // very long rebuilds, so cap generously and surface a timeout + // rather than hanging the whole recovery queue. + let cap = self.indexing_timeout() * 8; + let started = Instant::now(); + while self.is_indexing(alias) { + if started.elapsed() >= cap { + return Err(format!( + "rebuild for '{}' exceeded the {}s recovery cap", + alias, + cap.as_secs() + )); + } + tokio::time::sleep(Duration::from_secs(2)).await; + } + tracing::info!("🔧 Format recovery: rebuild complete for '{}'", alias); + Ok(()) + } + tui::ReindexLaunch::AlreadyRunning => { + // A rebuild is already in flight for this alias; wait it out. + // If it was a plain rebuild it may fail on the corrupt dir + // again — the next rebuild trigger re-detects and re-queues. + let cap = self.indexing_timeout() * 8; + let started = Instant::now(); + while self.is_indexing(alias) { + if started.elapsed() >= cap { + return Err(format!( + "in-flight rebuild for '{}' exceeded the {}s recovery cap", + alias, + cap.as_secs() + )); + } + tokio::time::sleep(Duration::from_secs(2)).await; + } + Ok(()) + } + tui::ReindexLaunch::Failed => Err(format!( + "could not start the recovery rebuild for '{}' (see log)", + alias + )), + } + } + /// Returns `true` if `alias` is currently (non-stale) indexing. /// /// Stale entries — those older than [`MAX_INDEXING_SECS`] — are lazily @@ -3670,14 +3860,29 @@ async fn trigger_symbol_rebuild( state.schedule_persist_repos_config(); } Ok(Err(e)) => { - tracing::error!("❌ Symbol rebuild failed for '{}': {}", alias_owned, e); + let msg = e.to_string(); state.end_indexing(&alias_owned); state .csharp_index_error - .insert(alias_owned.clone(), e.to_string()); - state - .csharp_index_status - .insert(alias_owned, CSharpIndexStatus::Error); + .insert(alias_owned.clone(), msg.clone()); + if ServeState::is_lmdb_format_corruption(&msg) { + tracing::warn!( + "⚠️ LMDB storage-format corruption for '{}' (data written by an older \ + storage-major) — queueing sequential wipe + full rebuild", + alias_owned + ); + // Recovery owns the outcome from here: show in-progress rather + // than Error; it flips to Ready on success or Error on failure. + state + .csharp_index_status + .insert(alias_owned.clone(), CSharpIndexStatus::Indexing); + state.enqueue_format_recovery(&alias_owned); + } else { + tracing::error!("❌ Symbol rebuild failed for '{}': {}", alias_owned, msg); + state + .csharp_index_status + .insert(alias_owned, CSharpIndexStatus::Error); + } } Err(e) => { tracing::error!( diff --git a/src/serve/tests.rs b/src/serve/tests.rs index d663e91b..80a4d8fe 100644 --- a/src/serve/tests.rs +++ b/src/serve/tests.rs @@ -2835,3 +2835,50 @@ async fn try_open_stores_honours_dimension_override_for_a_fresh_repo() { "a repo added with --model embeddinggemma-q4 must open at 768 dims, not the 384 default" ); } + +#[test] +fn is_lmdb_format_corruption_matches_known_lmdb_errors() { + let cases = [ + ( + "MDB_BAD_VALSIZE: Unsupported size of key/DB name/data, or wrong DUPFIXED size", + true, + ), + ("Symbol rebuild failed: heed -> MDB_BAD_VALSIZE", true), + ("storage error: wrong DUPFIXED size", true), + ("Unsupported size of key while opening vectordb", true), + ("scip-csharp failed with exit code 1", false), + ("MDB_NOTFOUND: No matching key/data pair found", false), + ("Task panicked: workspace load failed", false), + ("", false), + ]; + for (msg, expected) in cases { + assert_eq!( + ServeState::is_lmdb_format_corruption(msg), + expected, + "unexpected classification for {msg:?}" + ); + } +} + +#[tokio::test] +async fn enqueue_format_recovery_dedupes_and_starts_single_worker() { + let state = Arc::new(ServeState::new(ReposConfig::default(), None)); + // Same alias three times must collapse to one queued entry. The recovery + // worker is spawned but (current-thread test runtime) does not execute + // until an await point, so the queue content is asserted deterministically. + state.enqueue_format_recovery("ghost-repo"); + state.enqueue_format_recovery("ghost-repo"); + state.enqueue_format_recovery("ghost-repo"); + let len = state + .format_recovery_queue + .lock() + .expect("queue lock") + .len(); + assert_eq!(len, 1, "duplicate enqueues must collapse to one entry"); + assert!( + state + .format_recovery_worker_started + .load(std::sync::atomic::Ordering::Acquire), + "the first enqueue must start the recovery worker" + ); +} diff --git a/src/serve/tui.rs b/src/serve/tui.rs index b82834f5..9d31e60e 100644 --- a/src/serve/tui.rs +++ b/src/serve/tui.rs @@ -1060,7 +1060,7 @@ fn with_remote_stats(overlay: OverlayState, stats: RemoteStatsState) -> OverlayS /// Outcome of a TUI force-reindex launch — used to drive immediate footer /// feedback. Only describes whether the background task *started*; the actual /// indexing result is reported later via the status column / logs. -enum ReindexLaunch { +pub(crate) enum ReindexLaunch { /// Background reindex task spawned successfully. Started, /// A reindex was already running for this alias — request ignored. @@ -1071,7 +1071,7 @@ enum ReindexLaunch { /// Spawn a background force reindex task for the given repo alias. /// Follows the same flow as the HTTP `reindex_handler`. -fn spawn_force_reindex(alias: String, state: &Arc) -> ReindexLaunch { +pub(crate) fn spawn_force_reindex(alias: String, state: &Arc) -> ReindexLaunch { // Guard against concurrent reindex if !state.begin_indexing(&alias) { tracing::warn!(