From 5e6073baaa00b4fdfdf4f7e70847a70b1f41d0aa Mon Sep 17 00:00:00 2001 From: oxoxDev <164490987+oxoxDev@users.noreply.github.com> Date: Thu, 14 May 2026 08:37:55 +0530 Subject: [PATCH] fix(windows): retry-with-backoff for transient FS errors on auth-profiles.lock + .openhuman wipe (#9E, #9C, #4Y, #61, #5Q, #9F, #4M) (#1641) Co-authored-by: google-labs-jules[bot] <161369871+google-labs-jules[bot]@users.noreply.github.com> Co-authored-by: Claude Opus 4.7 (1M context) Co-authored-by: Steven Enamakel Co-authored-by: Zavian Wang <36817799+Zavianx@users.noreply.github.com> Co-authored-by: obchain <167975049+obchain@users.noreply.github.com> --- src/openhuman/credentials/profiles.rs | 76 ++++--- src/openhuman/memory/tree/read_rpc.rs | 302 ++++++++++++++------------ src/openhuman/util.rs | 271 ++++++++++++++++++++++- 3 files changed, 477 insertions(+), 172 deletions(-) diff --git a/src/openhuman/credentials/profiles.rs b/src/openhuman/credentials/profiles.rs index 0bf58ade3..cb4dd715f 100644 --- a/src/openhuman/credentials/profiles.rs +++ b/src/openhuman/credentials/profiles.rs @@ -469,11 +469,20 @@ impl AuthProfilesStore { let mut waited = 0_u64; let mut cleared_stale = false; loop { - match OpenOptions::new() - .create_new(true) - .write(true) - .open(&self.lock_path) - { + let open_result = crate::openhuman::util::retry_with_backoff( + "create auth profile lock", + 6, + 100, + || { + OpenOptions::new() + .create_new(true) + .write(true) + .open(&self.lock_path) + .context("open lock file") + }, + ); + + match open_result { Ok(mut file) => { // Issue #1612 — writing the pid line is what later lets // a future acquirer recognise a crashed owner; if the @@ -491,33 +500,38 @@ impl AuthProfilesStore { lock_path: self.lock_path.clone(), }); } - Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => { - // Issue #1612 — a previous openhuman crash can leave a - // stale auth-profiles.lock behind, after which every RPC - // path that touches the auth profile store fails for the - // 10s LOCK_TIMEOUT_MS window and the user gets stuck in a - // retry storm. Before falling back to the busy-wait, try - // once to peek at the writer's recorded PID and remove - // the lock if that process is no longer alive. Flag is - // flipped on the first probe (not only on success) so a - // live-pid / malformed / unreadable lock doesn't trigger - // a fresh sysinfo probe + log line on every busy-wait - // iteration. - if !cleared_stale { - cleared_stale = true; - if self.clear_lock_if_stale() { - continue; - } - } - if waited >= LOCK_TIMEOUT_MS { - anyhow::bail!("Timed out waiting for auth profile lock"); - } - thread::sleep(Duration::from_millis(LOCK_WAIT_MS)); - waited = waited.saturating_add(LOCK_WAIT_MS); - } Err(e) => { - return Err(e) - .with_context(|| "Failed to create auth profile lock".to_string()); + let is_already_exists = e + .chain() + .find_map(|e| e.downcast_ref::()) + .map_or(false, |ioe| ioe.kind() == std::io::ErrorKind::AlreadyExists); + + if is_already_exists { + // Issue #1612 — a previous openhuman crash can leave a + // stale auth-profiles.lock behind, after which every RPC + // path that touches the auth profile store fails for the + // 10s LOCK_TIMEOUT_MS window and the user gets stuck in a + // retry storm. Before falling back to the busy-wait, try + // once to peek at the writer's recorded PID and remove + // the lock if that process is no longer alive. Flag is + // flipped on the first probe (not only on success) so a + // live-pid / malformed / unreadable lock doesn't trigger + // a fresh sysinfo probe + log line on every busy-wait + // iteration. + if !cleared_stale { + cleared_stale = true; + if self.clear_lock_if_stale() { + continue; + } + } + if waited >= LOCK_TIMEOUT_MS { + anyhow::bail!("Timed out waiting for auth profile lock"); + } + thread::sleep(Duration::from_millis(LOCK_WAIT_MS)); + waited = waited.saturating_add(LOCK_WAIT_MS); + } else { + return Err(e).context("Failed to create auth profile lock"); + } } } } diff --git a/src/openhuman/memory/tree/read_rpc.rs b/src/openhuman/memory/tree/read_rpc.rs index a1d6c982f..81506956f 100644 --- a/src/openhuman/memory/tree/read_rpc.rs +++ b/src/openhuman/memory/tree/read_rpc.rs @@ -1254,7 +1254,7 @@ pub struct WipeAllResponse { /// can re-sync from scratch without leaving the app. pub async fn wipe_all_rpc(config: &Config) -> Result, String> { let cfg = config.clone(); - let resp = tokio::task::spawn_blocking(move || -> Result { + let (rows_deleted, sync_state_cleared) = tokio::task::spawn_blocking(move || -> Result<(u64, u64)> { // Tables to truncate. Order matters: `mem_tree_summaries` and // `mem_tree_buffers` both have `FOREIGN KEY (tree_id) REFERENCES // mem_tree_trees(id)` with `PRAGMA foreign_keys = ON`, so trees @@ -1283,45 +1283,11 @@ pub async fn wipe_all_rpc(config: &Config) -> Result Ok(total) })?; - // Filesystem cleanup. Each directory is best-effort: if one - // fails (permission denied, path doesn't exist) we keep going - // and report what we managed to remove. `email/` and the - // legacy bare `summaries/` are listed for back-compat — - // workspaces ingested before the raw-archive + wiki/ moves - // still have files there. Fresh installs only ever populate - // `raw/`, `wiki/`, `chat/`, and `document/`. - const DIRS: &[&str] = &["raw", "wiki", "chat", "document", "email", "summaries"]; - let content_root = cfg.memory_tree_content_root(); - let mut dirs_removed: Vec = Vec::new(); - for dir in DIRS { - let path = content_root.join(dir); - match std::fs::remove_dir_all(&path) { - Ok(()) => dirs_removed.push((*dir).to_string()), - Err(e) if e.kind() == std::io::ErrorKind::NotFound => {} - Err(e) => { - // Logical name (raw / wiki / chat / ...) is enough - // signal — the absolute path embeds the user's - // home directory. - log::warn!( - "[memory_tree::read::wipe] failed to remove dir={} err={e}", - dir - ); - } - } - } - // Composio sync-state lives in the unified memory store // (`/memory/memory.db`). Open it directly and // delete every key in the `composio-sync-state` namespace — // this clears each provider's `cursor` + `synced_ids` set so // the next sync re-fetches from the beginning. - // - // We do **not** swallow clear failures into `0`: callers (and - // the frontend `sync_state_cleared` contract) need to - // distinguish "nothing to clear" from "failed to clear, so - // the next sync may still be incremental." A missing DB is - // legitimately "nothing to clear"; a SQLite error is a - // failed wipe and propagates. let sync_state_cleared: u64 = { let unified_db = cfg.workspace_dir.join("memory").join("memory.db"); if !unified_db.exists() { @@ -1335,16 +1301,65 @@ pub async fn wipe_all_rpc(config: &Config) -> Result } }; - Ok(WipeAllResponse { - rows_deleted, - dirs_removed, - sync_state_cleared, - }) + Ok((rows_deleted, sync_state_cleared)) }) .await .map_err(|e| format!("wipe_all join error: {e}"))? .map_err(|e| format!("wipe_all: {e:#}"))?; + // Filesystem cleanup. Each directory is best-effort: if one + // fails (permission denied, path doesn't exist) we keep going + // and report what we managed to remove. `email/` and the + // legacy bare `summaries/` are listed for back-compat — + // workspaces ingested before the raw-archive + wiki/ moves + // still have files there. Fresh installs only ever populate + // `raw/`, `wiki/`, `chat/`, and `document/`. + // + // Use async retry to avoid blocking the executor during Windows sharing violations. + const DIRS: &[&str] = &["raw", "wiki", "chat", "document", "email", "summaries"]; + let content_root = config.memory_tree_content_root(); + let mut dirs_removed: Vec = Vec::new(); + for dir in DIRS { + let path = content_root.join(dir); + let remove_result = crate::openhuman::util::retry_with_backoff_async( + &format!("remove dir {}", dir), + 6, + 200, + || async { + tokio::fs::remove_dir_all(&path) + .await + .context("remove_dir_all") + }, + ) + .await; + + match remove_result { + Ok(()) => dirs_removed.push((*dir).to_string()), + Err(e) => { + let is_not_found = e + .chain() + .find_map(|e| e.downcast_ref::()) + .map_or(false, |ioe| ioe.kind() == std::io::ErrorKind::NotFound); + if !is_not_found { + // Logical name (raw / wiki / chat / ...) is enough + // signal — the absolute path embeds the user's + // home directory. + log::warn!( + "[memory_tree::read::wipe] failed to remove dir={} err={:#}", + dir, + e + ); + } + } + } + } + + let resp = WipeAllResponse { + rows_deleted, + dirs_removed, + sync_state_cleared, + }; + let log = format!( "memory_tree::read: wipe_all rows={} dirs={:?} sync_state={}", resp.rows_deleted, resp.dirs_removed, resp.sync_state_cleared @@ -1401,125 +1416,132 @@ pub struct ResetTreeResponse { /// inert fallback to a real Ollama model) and want to re-summarise /// existing data without paying the upstream sync cost again. /// -/// Three steps, each in its own SQL pass: +/// Three steps, executed in this order: /// 1. Truncate `mem_tree_summaries`, `mem_tree_trees`, /// `mem_tree_buffers`, `mem_tree_jobs`. The tree schema is /// derived state — chunks are the source of truth. -/// 2. Remove `/wiki/summaries/` on disk so stale -/// `.md` files don't drift from the SQL truth. -/// 3. Reset every chunk's `lifecycle_status` to +/// 2. Reset every chunk's `lifecycle_status` to /// `'pending_extraction'` and enqueue an `extract_chunk` job /// keyed on the chunk id. The async worker picks each up and /// re-runs entity extract → score → embed → append-to-buffer. /// Seals happen automatically as L0 buffers cross the gate. +/// 3. Remove `/wiki/summaries/` on disk so stale +/// `.md` files don't drift from the SQL truth. Done last (and +/// outside `spawn_blocking`) so the on-disk removal can use +/// async retry without blocking the worker thread. pub async fn reset_tree_rpc(config: &Config) -> Result, String> { use crate::openhuman::memory::tree::jobs::store as jobs_store; use crate::openhuman::memory::tree::jobs::types::{ExtractChunkPayload, NewJob}; let cfg = config.clone(); - let resp = tokio::task::spawn_blocking(move || -> Result { - // Step 1 — truncate tree state in one transaction. Chunks - // (`mem_tree_chunks`), the entity index, score rows, and the - // sync-state KV all stay intact. - // - // Order matters: `mem_tree_summaries` and `mem_tree_buffers` - // both have `FOREIGN KEY (tree_id) REFERENCES mem_tree_trees(id)`, - // and `PRAGMA foreign_keys = ON` is set. Trees must come last - // or SQLite throws "FOREIGN KEY constraint failed". `mem_tree_jobs` - // has no FK so its position is free. - // `mem_tree_entity_index` holds both leaf (chunk) and summary - // entity rows. Clearing it on reset prevents `top_entities` - // from counting orphan rows pointing at deleted summaries; - // the leaf rows get rebuilt naturally when the requeued - // `extract_chunk` jobs run for every chunk. - const TREE_TABLES: &[&str] = &[ - "mem_tree_summaries", - "mem_tree_buffers", - "mem_tree_jobs", - "mem_tree_entity_index", - "mem_tree_trees", - ]; - let tree_rows_deleted: u64 = with_connection(&cfg, |conn| { - let tx = conn.unchecked_transaction()?; - let mut total: u64 = 0; - for table in TREE_TABLES { - let n = tx - .execute(&format!("DELETE FROM {table}"), []) - .with_context(|| format!("delete from {table}"))?; - total += n as u64; - } - tx.commit()?; - Ok(total) - })?; - - // Step 2 — wipe the on-disk wiki/summaries tree. Best-effort: - // a missing folder is fine (fresh workspace). Other errors - // log + carry on — the SQL truth is what the rebuild relies on. - let summaries_dir = cfg - .memory_tree_content_root() - .join("wiki") - .join("summaries"); - match std::fs::remove_dir_all(&summaries_dir) { - Ok(()) => log::debug!("[memory_tree::read::reset_tree] removed wiki/summaries"), - Err(e) if e.kind() == std::io::ErrorKind::NotFound => {} - Err(e) => { - log::warn!("[memory_tree::read::reset_tree] failed to remove wiki/summaries: {e}") - } - } - - // Step 3 — flip every chunk back to `pending_extraction` and - // enqueue an `extract_chunk` job per id. Done in a single - // transaction so partial state is impossible: either the - // whole queue is in flight or nothing is. We use a chunked - // SELECT so very large workspaces don't materialise the - // entire id list in memory. - let (chunks_requeued, jobs_enqueued) = - with_connection(&cfg, |conn| -> anyhow::Result<(u64, u64)> { + let (tree_rows_deleted, chunks_requeued, jobs_enqueued) = + tokio::task::spawn_blocking(move || -> Result<(u64, u64, u64)> { + // Step 1 — truncate tree state in one transaction. + const TREE_TABLES: &[&str] = &[ + "mem_tree_summaries", + "mem_tree_buffers", + "mem_tree_jobs", + "mem_tree_entity_index", + "mem_tree_trees", + ]; + let tree_rows_deleted: u64 = with_connection(&cfg, |conn| { let tx = conn.unchecked_transaction()?; - let chunks_requeued = tx.execute( - "UPDATE mem_tree_chunks SET lifecycle_status = 'pending_extraction'", - [], - )? as u64; - let chunk_ids: Vec = { - let mut stmt = tx.prepare("SELECT id FROM mem_tree_chunks")?; - let rows = stmt - .query_map([], |r| r.get::<_, String>(0))? - .collect::>>() - .context("collect chunk ids")?; - rows - }; - let mut jobs_enqueued: u64 = 0; - for id in &chunk_ids { - let payload = ExtractChunkPayload { - chunk_id: id.clone(), - }; - let job = - NewJob::extract_chunk(&payload).context("build extract_chunk NewJob")?; - if jobs_store::enqueue_tx(&tx, &job) - .context("enqueue extract_chunk")? - .is_some() - { - jobs_enqueued += 1; - } + let mut total: u64 = 0; + for table in TREE_TABLES { + let n = tx + .execute(&format!("DELETE FROM {table}"), []) + .with_context(|| format!("delete from {table}"))?; + total += n as u64; } tx.commit()?; - Ok((chunks_requeued, jobs_enqueued)) + Ok(total) })?; - // Wake the worker pool so the freshly-enqueued jobs start - // running immediately rather than waiting for the next - // periodic poll. - crate::openhuman::memory::tree::jobs::wake_workers(); + // Step 2 — flip every chunk back to `pending_extraction` and + // enqueue an `extract_chunk` job per id. + let (chunks_requeued, jobs_enqueued) = + with_connection(&cfg, |conn| -> anyhow::Result<(u64, u64)> { + let tx = conn.unchecked_transaction()?; + let chunks_requeued = tx.execute( + "UPDATE mem_tree_chunks SET lifecycle_status = 'pending_extraction'", + [], + )? as u64; + let chunk_ids: Vec = { + let mut stmt = tx.prepare("SELECT id FROM mem_tree_chunks")?; + let rows = stmt + .query_map([], |r| r.get::<_, String>(0))? + .collect::>>() + .context("collect chunk ids")?; + rows + }; + let mut jobs_enqueued: u64 = 0; + for id in &chunk_ids { + let payload = ExtractChunkPayload { + chunk_id: id.clone(), + }; + let job = NewJob::extract_chunk(&payload) + .context("build extract_chunk NewJob")?; + if jobs_store::enqueue_tx(&tx, &job) + .context("enqueue extract_chunk")? + .is_some() + { + jobs_enqueued += 1; + } + } + tx.commit()?; + Ok((chunks_requeued, jobs_enqueued)) + })?; - Ok(ResetTreeResponse { - tree_rows_deleted, - chunks_requeued, - jobs_enqueued, + Ok((tree_rows_deleted, chunks_requeued, jobs_enqueued)) }) - }) - .await - .map_err(|e| format!("reset_tree join error: {e}"))? - .map_err(|e| format!("reset_tree: {e:#}"))?; + .await + .map_err(|e| format!("reset_tree join error: {e}"))? + .map_err(|e| format!("reset_tree: {e:#}"))?; + + // Step 3 — wipe the on-disk wiki/summaries tree. + // Use async retry to avoid blocking the executor during Windows sharing violations. + let summaries_dir = config + .memory_tree_content_root() + .join("wiki") + .join("summaries"); + let remove_result = crate::openhuman::util::retry_with_backoff_async( + "remove wiki/summaries", + 6, + 200, + || async { + tokio::fs::remove_dir_all(&summaries_dir) + .await + .context("remove_dir_all") + }, + ) + .await; + + match remove_result { + Ok(()) => log::debug!("[memory_tree::read::reset_tree] removed wiki/summaries"), + Err(e) => { + let is_not_found = e + .chain() + .find_map(|e| e.downcast_ref::()) + .map_or(false, |ioe| ioe.kind() == std::io::ErrorKind::NotFound); + if !is_not_found { + log::warn!( + "[memory_tree::read::reset_tree] failed to remove wiki/summaries: {:#}", + e + ) + } + } + } + + // Wake the worker pool. Done after the on-disk cleanup so jobs don't + // start racing against an in-progress directory removal; the small + // delay (at most the retry window on Windows) is acceptable. + crate::openhuman::memory::tree::jobs::wake_workers(); + + let resp = ResetTreeResponse { + tree_rows_deleted, + chunks_requeued, + jobs_enqueued, + }; let log = format!( "memory_tree::read: reset_tree tree_rows={} chunks={} jobs={}", diff --git a/src/openhuman/util.rs b/src/openhuman/util.rs index b11b0f656..3d6a2a25f 100644 --- a/src/openhuman/util.rs +++ b/src/openhuman/util.rs @@ -144,7 +144,7 @@ mod tests { #[test] fn test_truncate_cjk_characters() { // CJK characters (Chinese - each is 3 bytes) - let s = "这是一个测试消息用来触发崩溃的中文"; // 21 characters + let s = "这是一个测试消息用来触发崩溃 of the 中文"; // 21 characters let result = truncate_with_ellipsis(s, 16); assert!(result.ends_with("...")); assert!(result.is_char_boundary(result.len() - 1)); @@ -221,4 +221,273 @@ mod tests { assert_eq!(truncate_with_suffix(s, 5, "!!!"), "Hello!!!"); assert_eq!(truncate_with_suffix(s, 20, "!!!"), "Hello World"); } + + #[test] + fn test_retry_with_backoff_success_immediate() { + let mut calls = 0; + let result = retry_with_backoff("test", 3, 1, || { + calls += 1; + Ok::<_, anyhow::Error>(42) + }); + assert_eq!(result.unwrap(), 42); + assert_eq!(calls, 1); + } + + #[test] + fn test_retry_with_backoff_success_after_retries() { + let mut calls = 0; + let result = retry_with_backoff("test", 3, 1, || { + calls += 1; + if calls < 3 { + anyhow::bail!("__TEST_TRANSIENT__ error {}", calls); + } + Ok(42) + }); + assert_eq!(result.unwrap(), 42); + assert_eq!(calls, 3); + } + + #[tokio::test] + async fn test_retry_with_backoff_async_success_after_retries() { + use std::sync::atomic::{AtomicU32, Ordering}; + let calls = AtomicU32::new(0); + let result = retry_with_backoff_async("test_async", 3, 1, || async { + let c = calls.fetch_add(1, Ordering::SeqCst) + 1; + if c < 3 { + anyhow::bail!("__TEST_TRANSIENT__ error {}", c); + } + Ok(42) + }) + .await; + assert_eq!(result.unwrap(), 42); + assert_eq!(calls.load(Ordering::SeqCst), 3); + } + + #[test] + fn test_retry_with_backoff_failure_after_all_attempts() { + let mut calls = 0; + let result = retry_with_backoff("test", 3, 1, || { + calls += 1; + anyhow::bail!("__TEST_TRANSIENT__ error {}", calls); + #[allow(unreachable_code)] + Ok::(0) + }); + let err = result.unwrap_err(); + assert!(err.to_string().contains("test failed after 3 attempts")); + assert_eq!(calls, 3); + } + + #[test] + fn test_retry_with_backoff_bail_on_non_transient() { + let mut calls = 0; + let result = retry_with_backoff("test", 3, 1, || { + calls += 1; + anyhow::bail!("permanent error"); + #[allow(unreachable_code)] + Ok::(0) + }); + let err = result.unwrap_err(); + assert_eq!(err.to_string(), "permanent error"); + assert_eq!(calls, 1); + } + + #[tokio::test] + async fn test_retry_with_backoff_async_bail_on_non_transient() { + use std::sync::atomic::{AtomicU32, Ordering}; + let calls = AtomicU32::new(0); + let result = retry_with_backoff_async("test_async_bail", 3, 1, || async { + calls.fetch_add(1, Ordering::SeqCst); + anyhow::bail!("permanent error"); + #[allow(unreachable_code)] + Ok::(0) + }) + .await; + let err = result.unwrap_err(); + assert_eq!(err.to_string(), "permanent error"); + assert_eq!(calls.load(Ordering::SeqCst), 1); + } + + #[test] + fn test_retry_with_backoff_rejects_zero_attempts() { + let mut calls = 0; + let result = retry_with_backoff("zero_sync", 0, 1, || { + calls += 1; + Ok::(42) + }); + let err = result.unwrap_err(); + assert!( + err.to_string().contains("requires attempts > 0"), + "unexpected error message: {}", + err + ); + assert_eq!(calls, 0, "closure must not run when attempts == 0"); + } + + #[tokio::test] + async fn test_retry_with_backoff_async_rejects_zero_attempts() { + use std::sync::atomic::{AtomicU32, Ordering}; + let calls = AtomicU32::new(0); + let result = retry_with_backoff_async("zero_async", 0, 1, || async { + calls.fetch_add(1, Ordering::SeqCst); + Ok::(42) + }) + .await; + let err = result.unwrap_err(); + assert!( + err.to_string().contains("requires attempts > 0"), + "unexpected error message: {}", + err + ); + assert_eq!( + calls.load(Ordering::SeqCst), + 0, + "closure must not run when attempts == 0" + ); + } +} + +/// Helper to retry a filesystem operation with exponential backoff. +/// +/// Particularly useful on Windows where mandatory file locking often causes +/// transient `ERROR_SHARING_VIOLATION` (32) or `ERROR_ACCESS_DENIED` (5) +/// when multiple processes (or a stale handle) touch the same tree. +/// +/// Sleep `base_ms * 2^i` between attempts. Logs at `warn!` on retry and +/// `info!` on success-after-retry. +/// +/// **Note**: This is the synchronous version using `std::thread::sleep`. +/// Use `retry_with_backoff_async` in asynchronous contexts to avoid blocking +/// the executor. +pub fn retry_with_backoff( + op_name: &str, + attempts: u32, + base_ms: u64, + mut f: F, +) -> anyhow::Result +where + F: FnMut() -> anyhow::Result, +{ + anyhow::ensure!(attempts > 0, "{} requires attempts > 0", op_name); + + let mut last_err: Option = None; + + for i in 0..attempts { + match f() { + Ok(val) => { + if i > 0 { + tracing::info!(op = op_name, retries = i, "[util] succeeded after retries"); + } + return Ok(val); + } + Err(e) => { + if !is_transient_fs_error(&e) { + return Err(e); + } + + if i == attempts - 1 { + last_err = Some(e); + break; + } + + let sleep_ms = base_ms.saturating_mul(2u64.saturating_pow(i)).min(30_000); + tracing::warn!( + op = op_name, + attempt = i + 1, + max_attempts = attempts, + error = %e, + retry_in_ms = sleep_ms, + "[util] transient fs retry" + ); + + std::thread::sleep(std::time::Duration::from_millis(sleep_ms)); + } + } + } + + Err(last_err + .expect("attempts > 0") + .context(format!("{} failed after {} attempts", op_name, attempts))) +} + +/// Asynchronous version of `retry_with_backoff` using `tokio::time::sleep`. +pub async fn retry_with_backoff_async( + op_name: &str, + attempts: u32, + base_ms: u64, + mut f: F, +) -> anyhow::Result +where + F: FnMut() -> Fut, + Fut: std::future::Future>, +{ + anyhow::ensure!(attempts > 0, "{} requires attempts > 0", op_name); + + let mut last_err: Option = None; + + for i in 0..attempts { + match f().await { + Ok(val) => { + if i > 0 { + tracing::info!(op = op_name, retries = i, "[util] succeeded after retries"); + } + return Ok(val); + } + Err(e) => { + if !is_transient_fs_error(&e) { + return Err(e); + } + + if i == attempts - 1 { + last_err = Some(e); + break; + } + + let sleep_ms = base_ms.saturating_mul(2u64.saturating_pow(i)).min(30_000); + tracing::warn!( + op = op_name, + attempt = i + 1, + max_attempts = attempts, + error = %e, + retry_in_ms = sleep_ms, + "[util] transient fs retry" + ); + + tokio::time::sleep(std::time::Duration::from_millis(sleep_ms)).await; + } + } + } + + Err(last_err + .expect("attempts > 0") + .context(format!("{} failed after {} attempts", op_name, attempts))) +} + +/// Returns true if the error is a transient filesystem error that should be retried, +/// particularly on Windows where file locking is mandatory. +pub fn is_transient_fs_error(err: &anyhow::Error) -> bool { + // In tests, allow a specific error message to be treated as transient + // so we can verify the retry logic on all platforms. + if cfg!(test) && err.to_string().contains("__TEST_TRANSIENT__") { + return true; + } + + let io_err = err.chain().find_map(|e| e.downcast_ref::()); + + if let Some(io_err) = io_err { + #[cfg(windows)] + { + if let Some(code) = io_err.raw_os_error() { + // 5: ERROR_ACCESS_DENIED + // 32: ERROR_SHARING_VIOLATION + // 33: ERROR_LOCK_VIOLATION + // 1224: ERROR_USER_MAPPED_FILE + return code == 5 || code == 32 || code == 33 || code == 1224; + } + } + #[cfg(not(windows))] + { + let _ = io_err; + } + } + false }