From 9115002fa30073d280e21dae15dc031eb29362af Mon Sep 17 00:00:00 2001 From: Andrey Vasilevsky Date: Sat, 3 Oct 2026 21:59:28 -0700 Subject: [PATCH 1/2] fix(storage): wait out a busy graph DB instead of discarding it A session that started while another process had graph.db open (any sibling session or daemon mid-load or mid-persist) treated the lock error as corruption: it set the healthy DB aside as graph.db.corrupt and redirected every project to a fresh, empty generation. On Linux the stale-LOCK cleanup could also delete a live holder's LOCK file, because its flock probe cannot see RocksDB's fcntl lock, letting two processes open the same DB. - RocksDBBackend::open reports lock contention as GraphError::Locked. - open_waiting_for_lock replaces open_with_stale_lock_recovery: it retries with backoff for a bounded time and never removes LOCK. The OS drops a lock when its holder exits, so a leftover LOCK file never blocked an open; removing one only ever risked a second writer. fs2 is no longer needed. - Every open of the shared DB goes through memory::open_shared_graph_db (10s wait), including vector claim/save/load and cross-project impact. - open_persistent_graph redirects only on a real load failure, never on Locked. - Saved file hashes are ignored when the project's graph loads empty. After any redirect they made the indexer skip every unchanged file, leaving the project with no symbols until a file changed. Verified end to end with the release binary (live holder for 4s and 20s, SIGKILLed holder, damaged CURRENT) and by new crate and integration tests on macOS, SLES 15-SP4 and Windows. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017rVbt7rENTwXkdHt3Bpgb5 --- Cargo.lock | 11 - Cargo.toml | 2 - .../codegraph-server/src/ai_query/engine.rs | 16 +- crates/codegraph-server/src/backend.rs | 20 +- crates/codegraph-server/src/domain/impact.rs | 8 +- crates/codegraph-server/src/main.rs | 2 +- crates/codegraph-server/src/mcp/server.rs | 41 ++-- crates/codegraph-server/src/memory.rs | 13 ++ .../tests/graph_db_recovery_test.rs | 142 ++++++++++++ crates/codegraph/Cargo.toml | 4 +- crates/codegraph/src/error.rs | 13 ++ .../codegraph/src/storage/rocksdb_backend.rs | 204 ++++++++---------- 12 files changed, 306 insertions(+), 170 deletions(-) create mode 100644 crates/codegraph-server/tests/graph_db_recovery_test.rs diff --git a/Cargo.lock b/Cargo.lock index 98c39e9..2bd6e00 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -534,7 +534,6 @@ name = "codegraph" version = "0.2.0" dependencies = [ "criterion", - "fs2", "log", "rocksdb", "serde", @@ -1734,16 +1733,6 @@ dependencies = [ "percent-encoding", ] -[[package]] -name = "fs2" -version = "0.4.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9564fc758e15025b46aa6643b1b77d047d1a56a1aea6e01002ac0c7026876213" -dependencies = [ - "libc", - "winapi", -] - [[package]] name = "fsevent-sys" version = "4.1.0" diff --git a/Cargo.toml b/Cargo.toml index 9b0bcae..4ced4be 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -74,8 +74,6 @@ thiserror = "1.0" rocksdb = "0.22" # Pin libz-sys to 1.1.25 — v1.1.26 has broken vendored zlib build on macOS/Windows libz-sys = "=1.1.25" -# Cross-platform advisory file locks (Linux fcntl / Windows LockFileEx) -fs2 = "0.4" log = "0.4" uuid = { version = "1.0", features = ["v4", "serde"] } diff --git a/crates/codegraph-server/src/ai_query/engine.rs b/crates/codegraph-server/src/ai_query/engine.rs index f709506..ec3e061 100644 --- a/crates/codegraph-server/src/ai_query/engine.rs +++ b/crates/codegraph-server/src/ai_query/engine.rs @@ -971,15 +971,13 @@ impl QueryEngine { /// rebuild that follows are loadable if it is interrupted. Returns whether /// the project was taken - see [`claim_project`]. fn claim_vector_set(slug: &str, stamp: &str) -> std::result::Result { - use codegraph::RocksDBBackend; - let db_path = crate::memory::shared_graph_db_path().map_err(|e| format!("{e}"))?; if let Some(parent) = db_path.parent() { std::fs::create_dir_all(parent) .map_err(|e| format!("Failed to create ~/.codegraph: {e}"))?; } - let rocks = - RocksDBBackend::open(&db_path).map_err(|e| format!("Failed to open graph.db: {e}"))?; + let rocks = crate::memory::open_shared_graph_db(&db_path) + .map_err(|e| format!("Failed to open graph.db: {e}"))?; let mut namespaced = NamespacedBackend::new(Box::new(rocks), slug); claim_project(&mut namespaced, slug, stamp) @@ -1013,8 +1011,6 @@ impl QueryEngine { complete: bool, stamp: &str, ) -> std::result::Result { - use codegraph::RocksDBBackend; - if vecs.is_empty() { return Ok(false); } @@ -1025,8 +1021,8 @@ impl QueryEngine { .map_err(|e| format!("Failed to create ~/.codegraph: {e}"))?; } - let rocks = - RocksDBBackend::open(&db_path).map_err(|e| format!("Failed to open graph.db: {e}"))?; + let rocks = crate::memory::open_shared_graph_db(&db_path) + .map_err(|e| format!("Failed to open graph.db: {e}"))?; let mut namespaced = NamespacedBackend::new(Box::new(rocks), slug); if !store_vectors(&mut namespaced, vecs, stamp)? { @@ -1058,8 +1054,6 @@ impl QueryEngine { /// engine is configured to build; a set built from other text is left /// untouched for whoever wrote it. Returns the number of vectors loaded. pub async fn load_symbol_vectors(&self, slug: &str) -> VectorLoad { - use codegraph::RocksDBBackend; - let Some(stamp) = self.embed_stamp().await else { tracing::warn!("[QueryEngine] No vector engine attached - cannot load symbol vectors"); return VectorLoad::Unreadable; @@ -1077,7 +1071,7 @@ impl QueryEngine { return VectorLoad::Absent; } - let rocks = match RocksDBBackend::open(&db_path) { + let rocks = match crate::memory::open_shared_graph_db(&db_path) { Ok(r) => r, Err(e) => { tracing::warn!("[QueryEngine] Failed to open graph.db for vectors: {}", e); diff --git a/crates/codegraph-server/src/backend.rs b/crates/codegraph-server/src/backend.rs index 88c54ce..94e23b9 100644 --- a/crates/codegraph-server/src/backend.rs +++ b/crates/codegraph-server/src/backend.rs @@ -515,7 +515,7 @@ impl CodeGraphBackend { workspace: &std::path::Path, graph: &codegraph::CodeGraph, ) -> std::result::Result<(), String> { - use codegraph::{NamespacedBackend, RocksDBBackend, StorageBackend}; + use codegraph::{NamespacedBackend, StorageBackend}; // Ephemeral workspaces (test harness tempdirs) skip // persistence to the shared graph.db entirely. The in-memory @@ -532,7 +532,7 @@ impl CodeGraphBackend { .map_err(|e| format!("Failed to create ~/.codegraph: {e}"))?; } - let mut rocks = RocksDBBackend::open_with_stale_lock_recovery(&db_path) + let mut rocks = crate::memory::open_shared_graph_db(&db_path) .map_err(|e| format!("Failed to open graph.db: {e}"))?; let registry_value = serde_json::json!({ @@ -1206,16 +1206,15 @@ impl LanguageServer for CodeGraphBackend { } Err(e) => { tracing::error!( - "LSP: RocksDB graph.db open failed: {e} — running in-memory only \ - this session. Changes will NOT persist across restarts." + "LSP: could not load the persisted graph ({e}) - starting without it." ); self.client .log_message( MessageType::ERROR, format!( - "CodeGraph: graph database open failed ({e}). \ - Index is volatile this session — restart to retry; \ - if the error persists, check ~/.codegraph/graph.db for a stale LOCK file." + "CodeGraph: could not load the saved index ({e}), so this session \ + starts without it. If another CodeGraph process keeps \ + ~/.codegraph/graph.db busy, restart once it is idle." ), ) .await; @@ -1223,6 +1222,13 @@ impl LanguageServer for CodeGraphBackend { } } + // The saved hashes (loaded in `initialize`) vouch that files are already + // in the graph. Without a persisted graph they would make the indexer + // skip every unchanged file and leave the session empty. + if !loaded_from_persistence { + self.index_state.lock().await.clear(); + } + // Run incremental indexing: hash-based dedup skips unchanged files. // - Fresh start (no persisted graph): full index if indexOnStartup=true // - Loaded from persistence: incremental pass to catch changed files diff --git a/crates/codegraph-server/src/domain/impact.rs b/crates/codegraph-server/src/domain/impact.rs index 3d53841..2086b97 100644 --- a/crates/codegraph-server/src/domain/impact.rs +++ b/crates/codegraph-server/src/domain/impact.rs @@ -7,9 +7,7 @@ use crate::ai_query::QueryEngine; use crate::domain::node_props; -use codegraph::{ - CodeGraph, Direction, EdgeType, NamespacedBackend, NodeId, RocksDBBackend, StorageBackend, -}; +use codegraph::{CodeGraph, Direction, EdgeType, NamespacedBackend, NodeId, StorageBackend}; use serde::Serialize; use std::collections::HashSet; use tokio::sync::RwLock; @@ -341,7 +339,7 @@ fn find_cross_project_consumers( // Open RocksDB, scan registry, then DROP the connection before per-project // loading — RocksDB uses exclusive locks, so only one connection at a time. let entries = { - let rocks = match RocksDBBackend::open(&db_path) { + let rocks = match crate::memory::open_shared_graph_db(&db_path) { Ok(r) => r, Err(e) => { tracing::warn!("[cross-project] Failed to open graph.db: {}", e); @@ -401,7 +399,7 @@ fn find_cross_project_consumers( }) .unwrap_or_else(|| slug.clone()); - let other_rocks = match RocksDBBackend::open(&db_path) { + let other_rocks = match crate::memory::open_shared_graph_db(&db_path) { Ok(r) => r, Err(_) => continue, }; diff --git a/crates/codegraph-server/src/main.rs b/crates/codegraph-server/src/main.rs index 1418a4e..930f90e 100644 --- a/crates/codegraph-server/src/main.rs +++ b/crates/codegraph-server/src/main.rs @@ -276,7 +276,7 @@ fn classify_panic(payload: &str, location: &str) -> (&'static str, &'static str) /// Strategy: panic / SIGINT / SIGTERM all funnel into `process::exit`. /// At process exit the kernel releases all fcntl / LockFileEx grants, /// so the next launch sees only the `LOCK` *file* (no live holder), -/// which `RocksDBBackend::open_with_stale_lock_recovery` clears. WAL +/// which RocksDB reuses without complaint. WAL /// durability is per-write, so any in-flight batch is either fully /// applied or fully discarded on next open — `exit` skipping `Drop` is /// a safe tradeoff here. diff --git a/crates/codegraph-server/src/mcp/server.rs b/crates/codegraph-server/src/mcp/server.rs index fa8abc2..6badedc 100644 --- a/crates/codegraph-server/src/mcp/server.rs +++ b/crates/codegraph-server/src/mcp/server.rs @@ -92,7 +92,7 @@ use crate::index_state::IndexState; use crate::indexer::{IndexConfig, Indexer}; use crate::memory::{self, MemoryManager}; use crate::parser_registry::ParserRegistry; -use codegraph::{CodeGraph, NamespacedBackend, RocksDBBackend, StorageBackend}; +use codegraph::{CodeGraph, NamespacedBackend, StorageBackend}; use serde_json::Value; use std::path::PathBuf; use std::sync::Arc; @@ -200,9 +200,9 @@ impl McpBackend { } Err(e) => { tracing::error!( - "RocksDB graph.db open failed: {e} — running in-memory only this session. \ - Changes will NOT persist across restarts. Inspect ~/.codegraph/graph.db and \ - ensure no other codegraph-server process is running." + "Could not load the persisted graph ({e}) - this session re-indexes the \ + workspace instead. If another codegraph process keeps ~/.codegraph/graph.db \ + busy, restart this session once it is idle." ); Arc::new(RwLock::new( CodeGraph::in_memory().expect("Failed to create in-memory graph"), @@ -352,6 +352,10 @@ impl McpBackend { match result { Ok(graph) => Ok(graph), + // Another process kept the DB open past the wait. That is + // contention, not damage: redirecting here would abandon every + // project's graph because a sibling session was mid-persist. + Err(e @ codegraph::GraphError::Locked { .. }) => Err(format!("graph.db: {e}")), Err(e) => { // RocksDB reported corruption gracefully (didn't AV). // Redirect to a fresh generation and retry once so the @@ -363,6 +367,7 @@ impl McpBackend { Self::sweep_stale_graph_dbs(parent, &fresh); } Self::load_persistent_graph_inner(&fresh, slug) + .map_err(|e| format!("graph.db load failed: {e}")) } } } @@ -519,20 +524,13 @@ impl McpBackend { fn load_persistent_graph_inner( db_path: &std::path::Path, slug: &str, - ) -> Result { - // Stale-LOCK recovery: a prior crash can leave LOCK in place. The - // recovery variant only clobbers it after probing for a live holder — - // a healthy concurrent process is still respected. - let rocks = RocksDBBackend::open_with_stale_lock_recovery(db_path) - .map_err(|e| format!("Failed to open graph.db: {e}"))?; + ) -> Result { + let rocks = memory::open_shared_graph_db(db_path)?; let namespaced = NamespacedBackend::new(Box::new(rocks), slug); - let mut graph = CodeGraph::with_backend(Box::new(namespaced)) - .map_err(|e| format!("Failed to load graph: {e}"))?; + let mut graph = CodeGraph::with_backend(Box::new(namespaced))?; // Detach to release the RocksDB lock — all data is now in memory - graph - .detach_storage() - .map_err(|e| format!("Failed to detach storage: {e}"))?; + graph.detach_storage()?; Ok(graph) } @@ -580,7 +578,7 @@ impl McpBackend { .map_err(|e| format!("Failed to create ~/.codegraph: {e}"))?; } - let mut rocks = RocksDBBackend::open_with_stale_lock_recovery(&db_path) + let mut rocks = memory::open_shared_graph_db(&db_path) .map_err(|e| format!("Failed to open graph.db for persist: {e}"))?; // Write project registry entry (un-namespaced, global key) @@ -633,7 +631,7 @@ impl McpBackend { return Ok(vec![]); } - let rocks = RocksDBBackend::open_with_stale_lock_recovery(&db_path) + let rocks = memory::open_shared_graph_db(&db_path) .map_err(|e| format!("Failed to open graph.db: {e}"))?; let entries = rocks @@ -903,7 +901,16 @@ impl McpBackend { } /// Load saved file hashes from disk. Returns true if state was loaded. + /// + /// The hashes vouch that a file's symbols are already in the graph, so + /// they only load alongside a graph that has some. Against an empty graph + /// (a redirected or unloadable DB) they would make the indexer skip every + /// unchanged file and leave the session with no symbols at all. pub async fn load_index_state(&self) -> bool { + if self.graph.read().await.node_count() == 0 { + tracing::info!("No persisted graph to resume from - indexing every file"); + return false; + } let mut state = self.index_state.lock().await; let count = state.load(); count > 0 diff --git a/crates/codegraph-server/src/memory.rs b/crates/codegraph-server/src/memory.rs index fdbe3ad..e9ae983 100644 --- a/crates/codegraph-server/src/memory.rs +++ b/crates/codegraph-server/src/memory.rs @@ -105,6 +105,19 @@ pub(crate) fn shared_graph_db_path() -> Result { Ok(graph_db_path_for_generation(&dir, graph_db_generation())) } +/// How long an open of the shared graph DB waits for another process to let go +/// of it. Sessions and daemons hold it only for a load or a persist. +pub(crate) const SHARED_GRAPH_DB_LOCK_WAIT: std::time::Duration = + std::time::Duration::from_secs(10); + +/// Open the shared graph DB at `path`, waiting out another process that has it +/// open. Every open of the shared DB goes through here: RocksDB admits one +/// handle at a time, so an open that gives up at the first refusal fails +/// whenever a sibling session happens to be persisting. +pub(crate) fn open_shared_graph_db(path: &Path) -> codegraph::Result { + codegraph::RocksDBBackend::open_waiting_for_lock(path, SHARED_GRAPH_DB_LOCK_WAIT) +} + /// Redirect the shared graph DB to a brand-new directory by bumping the /// generation pointer; returns the new path. The poisoned old directory is /// NOT touched here — renaming/deleting it is best-effort cleanup that can diff --git a/crates/codegraph-server/tests/graph_db_recovery_test.rs b/crates/codegraph-server/tests/graph_db_recovery_test.rs new file mode 100644 index 0000000..6cf4336 --- /dev/null +++ b/crates/codegraph-server/tests/graph_db_recovery_test.rs @@ -0,0 +1,142 @@ +// Copyright 2024-2026 Andrey Vasilevsky +// SPDX-License-Identifier: Apache-2.0 + +//! End-to-end checks of how a session treats the shared graph DB when it can't +//! simply open it: another process holding it, or a damaged DB. +//! +//! Each test runs the real `codegraph-server --mcp` binary against a temporary +//! workspace and a temporary HOME, so the user's own `~/.codegraph` is never +//! touched. + +use codegraph::RocksDBBackend; +use serde_json::{json, Value}; +use std::io::{BufRead, BufReader, Write}; +use std::path::{Path, PathBuf}; +use std::process::{Command, Stdio}; +use std::time::Duration; + +const MARKER: &str = "graph_db_recovery_marker_fn"; + +struct Env { + _root: tempfile::TempDir, + workspace: PathBuf, + home: PathBuf, +} + +impl Env { + fn new() -> Self { + let root = tempfile::tempdir().unwrap(); + let workspace = root.path().join("ws"); + let home = root.path().join("home"); + std::fs::create_dir_all(&workspace).unwrap(); + std::fs::create_dir_all(&home).unwrap(); + std::fs::write(workspace.join("lib.rs"), format!("fn {MARKER}() {{}}\n")).unwrap(); + Self { + _root: root, + workspace, + home, + } + } + + fn codegraph_dir(&self) -> PathBuf { + self.home.join(".codegraph") + } + + /// The DB a session opens when no redirect has happened. + fn first_generation_db(&self) -> PathBuf { + self.codegraph_dir().join("graph.db") + } + + fn redirected(&self) -> bool { + self.codegraph_dir().join("graph.generation").exists() + } + + /// Run one MCP session and report whether a symbol search finds the marker. + fn session_finds_marker(&self) -> bool { + let mut child = Command::new(env!("CARGO_BIN_EXE_codegraph-server")) + .args(["--mcp", "--graph-only", "--workspace"]) + .arg(&self.workspace) + .env("HOME", &self.home) + .env("USERPROFILE", &self.home) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .spawn() + .unwrap(); + let mut stdin = child.stdin.take().unwrap(); + let mut stdout = BufReader::new(child.stdout.take().unwrap()); + let mut request = |id: u64, method: &str, params: Value| -> Value { + let msg = json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params}); + writeln!(stdin, "{msg}").unwrap(); + loop { + let mut line = String::new(); + assert!(stdout.read_line(&mut line).unwrap() > 0, "server exited"); + let reply: Value = serde_json::from_str(&line).unwrap(); + if reply["id"] == id { + return reply; + } + } + }; + request( + 1, + "initialize", + json!({"protocolVersion": "2024-11-05", "capabilities": {}, + "clientInfo": {"name": "test", "version": "1"}}), + ); + let reply = request( + 2, + "tools/call", + json!({"name": "codegraph_symbol_search", "arguments": {"query": MARKER}}), + ); + drop(request); + drop(stdin); + assert!(child.wait().unwrap().success()); + reply.to_string().contains(MARKER) + } +} + +/// Hold the DB open from this process for `hold`, as a sibling session does +/// while it loads or persists. +fn hold_db(path: &Path, hold: Duration) -> std::thread::JoinHandle<()> { + let holder = RocksDBBackend::open(path).unwrap(); + std::thread::spawn(move || { + std::thread::sleep(hold); + drop(holder); + }) +} + +#[test] +fn a_session_waits_for_a_sibling_holding_the_db_instead_of_abandoning_it() { + let env = Env::new(); + assert!(env.session_finds_marker(), "seed session"); + + let holder = hold_db(&env.first_generation_db(), Duration::from_secs(2)); + assert!( + env.session_finds_marker(), + "session started while the DB was held" + ); + holder.join().unwrap(); + + assert!( + !env.redirected(), + "lock contention must not redirect the DB" + ); + assert!(env.session_finds_marker(), "later session"); +} + +#[test] +fn a_session_after_a_redirect_reindexes_files_the_old_db_held() { + let env = Env::new(); + assert!(env.session_finds_marker(), "seed session"); + + // An unreadable DB: the session redirects to a fresh, empty generation. + std::fs::write( + env.first_generation_db().join("CURRENT"), + "MANIFEST-999999\n", + ) + .unwrap(); + assert!(env.session_finds_marker(), "session that redirected"); + assert!(env.redirected()); + + assert!(env.session_finds_marker(), "session after the redirect"); +} diff --git a/crates/codegraph/Cargo.toml b/crates/codegraph/Cargo.toml index 51f5900..bd66548 100644 --- a/crates/codegraph/Cargo.toml +++ b/crates/codegraph/Cargo.toml @@ -15,13 +15,11 @@ categories = ["database-implementations", "data-structures", "development-tools" [features] default = ["rocksdb-backend"] -rocksdb-backend = ["dep:rocksdb", "dep:fs2"] +rocksdb-backend = ["dep:rocksdb"] [dependencies] # Storage backend rocksdb = { workspace = true, optional = true } -# Stale-LOCK detection for crash recovery -fs2 = { workspace = true, optional = true } # Serialization serde.workspace = true diff --git a/crates/codegraph/src/error.rs b/crates/codegraph/src/error.rs index e382ff1..4e47e0a 100644 --- a/crates/codegraph/src/error.rs +++ b/crates/codegraph/src/error.rs @@ -26,6 +26,19 @@ pub enum GraphError { source: Option>, }, + /// The database is held by another open handle, usually another process. + /// + /// Contention, not damage: the database itself is fine and opens once the + /// holder lets go. + #[error("Database at {path:?} is in use by another process")] + Locked { + /// Database directory + path: PathBuf, + /// The storage engine's lock error + #[source] + source: Box, + }, + /// Node not found in the graph #[error("Node not found: {node_id}")] NodeNotFound { diff --git a/crates/codegraph/src/storage/rocksdb_backend.rs b/crates/codegraph/src/storage/rocksdb_backend.rs index eede75a..fdbd64f 100644 --- a/crates/codegraph/src/storage/rocksdb_backend.rs +++ b/crates/codegraph/src/storage/rocksdb_backend.rs @@ -11,6 +11,7 @@ use crate::error::{GraphError, Result}; use rocksdb::{Options, WriteBatch, DB}; use std::path::Path; use std::sync::Arc; +use std::time::{Duration, Instant}; /// RocksDB-backed persistent storage. /// @@ -47,41 +48,41 @@ impl RocksDBBackend { // open-time option can catch. opts.set_wal_recovery_mode(rocksdb::DBRecoveryMode::PointInTime); - let db = DB::open(&opts, path.as_ref()).map_err(|e| { - GraphError::storage( - format!("Failed to open RocksDB at {:?}", path.as_ref()), - Some(e), - ) - })?; + let db = DB::open(&opts, path.as_ref()).map_err(|e| open_error(path.as_ref(), e))?; Ok(Self { db: Arc::new(db) }) } - /// Open with recovery from a stale `LOCK` file left by a prior crash. + /// Open, waiting up to `max_wait` for another open handle to release the + /// database. + /// + /// RocksDB allows one open handle per database, and the holder is usually + /// another process that keeps it open only for a load or a persist, so + /// contention clears within moments. This retries with backoff while + /// [`Self::open`] reports [`GraphError::Locked`], and returns that error + /// once `max_wait` has elapsed. + /// + /// It never removes the `LOCK` file. The OS releases a lock when its + /// holder exits, crashed or not, so a `LOCK` file left behind does not + /// block the next open, and removing one that is still held would let a + /// second process open the same database. /// - /// Falls back to a single retry only when the original failure looks - /// lock-related AND an advisory-lock probe of `/LOCK` succeeds — - /// the probe is the authoritative signal that no live process still - /// holds the inode. Without that double check we would happily steal a - /// lock from a healthy concurrent process. + /// # Errors /// - /// Use this for production open-paths (server startup, persist). Tests - /// and tools that want strict open-time conflict detection should - /// continue to call [`Self::open`]. - pub fn open_with_stale_lock_recovery>(path: P) -> Result { - let path_ref = path.as_ref(); - match Self::open(path_ref) { - Ok(b) => Ok(b), - Err(e) => { - if is_lock_error(&e) && try_clear_stale_lock(path_ref) { - log::warn!( - "RocksDB at {:?} had a stale LOCK from a prior crash; cleared and retrying", - path_ref, + /// Returns [`GraphError::Locked`] if the database is still held after + /// `max_wait`, or [`GraphError::Storage`] for any other open failure. + pub fn open_waiting_for_lock>(path: P, max_wait: Duration) -> Result { + let deadline = Instant::now() + max_wait; + let mut delay = Duration::from_millis(20); + loop { + match Self::open(path.as_ref()) { + Err(GraphError::Locked { .. }) if Instant::now() < deadline => { + std::thread::sleep( + delay.min(deadline.saturating_duration_since(Instant::now())), ); - Self::open(path_ref) - } else { - Err(e) + delay = (delay * 2).min(Duration::from_millis(500)); } + result => return result, } } } @@ -168,62 +169,28 @@ impl RocksDBBackend { } } -/// Heuristic: does this storage error look like a `LOCK`-file failure? +/// Classify a failed `DB::open`: lock contention becomes [`GraphError::Locked`], +/// anything else [`GraphError::Storage`]. /// -/// RocksDB's `Error` type is opaque (string-only), so substring matching -/// is the only portable detector. Patterns checked here are the literal -/// strings the underlying C++ layer emits across the platforms we ship -/// (`IOError: While lock file ... LOCK`, `Resource temporarily unavailable`, -/// `lock hold`). False positives are safe — they only cause a probe; the -/// probe itself is what authorises cleanup. -fn is_lock_error(e: &GraphError) -> bool { - use std::error::Error; - // Walk the source chain so we catch the underlying rocksdb::Error too. - let mut s = format!("{e}"); - let mut src: Option<&(dyn Error + 'static)> = e.source(); - while let Some(inner) = src { - s.push('\n'); - s.push_str(&inner.to_string()); - src = inner.source(); - } - let needles = [ - "lock", - "LOCK", - "Resource temporarily unavailable", - "lock hold", +/// RocksDB's error is string-only. These are the messages its `LockFile` +/// returns when another handle holds the database: another process on POSIX, +/// this process on POSIX, and any holder on Windows (the `LOCK` file is opened +/// without sharing there). +fn open_error(path: &Path, e: rocksdb::Error) -> GraphError { + const LOCK_HELD: [&str; 3] = [ + "While lock file", + "lock hold by current process", + "Failed to create lock file", ]; - needles.iter().any(|n| s.contains(n)) -} - -/// Probe `/LOCK` for a live holder; remove it if none. -/// -/// Returns `true` only when (a) the file exists, (b) an advisory-lock -/// probe succeeds — meaning no other process holds an exclusive lock on -/// the inode — and (c) the file was successfully removed. Any other -/// outcome returns `false` so the caller surfaces the original error -/// (real conflict, permission issue, missing parent dir, etc.). -/// -/// The advisory probe uses the same lock primitive RocksDB itself uses -/// (fcntl on POSIX, LockFileEx on Windows), so a healthy concurrent -/// process is reliably detected and not stomped. -fn try_clear_stale_lock(db_path: &Path) -> bool { - use fs2::FileExt; - use std::fs::OpenOptions; - - let lock_path = db_path.join("LOCK"); - if !lock_path.exists() { - return false; - } - let file = match OpenOptions::new().read(true).write(true).open(&lock_path) { - Ok(f) => f, - Err(_) => return false, - }; - if file.try_lock_exclusive().is_err() { - return false; + let message = e.to_string(); + if LOCK_HELD.iter().any(|m| message.contains(m)) { + GraphError::Locked { + path: path.to_path_buf(), + source: Box::new(e), + } + } else { + GraphError::storage(format!("Failed to open RocksDB at {path:?}"), Some(e)) } - let _ = FileExt::unlock(&file); - drop(file); - std::fs::remove_file(&lock_path).is_ok() } impl StorageBackend for RocksDBBackend { @@ -467,50 +434,61 @@ mod tests { } #[test] - fn test_stale_lock_recovery_clears_orphaned_lock() { - // Simulate the post-crash state: a LOCK file exists but no - // process holds an advisory lock on it. open() returns Err - // (because we manually fcntl-locked it from a side-channel that - // got closed); open_with_stale_lock_recovery() should clear and - // succeed. + fn test_leftover_lock_file_does_not_block_open() { + // A holder that exits, crashed or not, leaves its LOCK file behind + // but not its lock: the next open succeeds and the file is reused. let temp_dir = TempDir::new().unwrap(); let db_path = temp_dir.path().to_path_buf(); - std::fs::create_dir_all(&db_path).unwrap(); + drop(RocksDBBackend::open(&db_path).unwrap()); + assert!(db_path.join("LOCK").exists()); - // First clean open creates the LOCK as a side-effect of RocksDB - // initialising the directory. - { - let backend = RocksDBBackend::open(&db_path).unwrap(); - drop(backend); - } + let backend = RocksDBBackend::open(&db_path).unwrap(); + backend.get(b"anything").unwrap(); + } - // Recreate a LOCK file by hand — emulates a leftover from a - // killed process where the kernel released the fcntl lock but - // never removed the file (Windows-shaped state). - let lock_path = db_path.join("LOCK"); - if !lock_path.exists() { - std::fs::write(&lock_path, b"").unwrap(); - } + #[test] + fn test_open_reports_a_held_database_as_locked() { + let temp_dir = TempDir::new().unwrap(); + let db_path = temp_dir.path().to_path_buf(); + let _holder = RocksDBBackend::open(&db_path).unwrap(); - // No one holds an advisory lock on it → recovery should succeed. - let backend = RocksDBBackend::open_with_stale_lock_recovery(&db_path).unwrap(); - backend.get(b"anything").unwrap(); + let err = RocksDBBackend::open(&db_path).err().unwrap(); + assert!(matches!(err, GraphError::Locked { .. }), "got {err}"); } #[test] - fn test_stale_lock_recovery_with_live_holder_does_not_deadlock_or_panic() { - // A live holder is alive in the same process. The recovery API - // must not deadlock, panic, or corrupt state. Either Ok or Err - // is acceptable as the return — Windows refuses the LOCK probe - // outright (sharing violation), POSIX may permit it under - // same-process fcntl semantics. The load-bearing assertion is - // simply that the call returns within a reasonable time. + fn test_waiting_open_gives_up_after_max_wait_and_keeps_the_lock_file() { let temp_dir = TempDir::new().unwrap(); let db_path = temp_dir.path().to_path_buf(); - std::fs::create_dir_all(&db_path).unwrap(); + let mut holder = RocksDBBackend::open(&db_path).unwrap(); + holder.put(b"k", b"v").unwrap(); + + let started = std::time::Instant::now(); + let err = RocksDBBackend::open_waiting_for_lock(&db_path, Duration::from_millis(300)) + .err() + .unwrap(); + assert!(matches!(err, GraphError::Locked { .. }), "got {err}"); + assert!(started.elapsed() >= Duration::from_millis(300)); + assert!(db_path.join("LOCK").exists()); + // The holder is undisturbed. + assert_eq!(holder.get(b"k").unwrap(), Some(b"v".to_vec())); + } - let _holder = RocksDBBackend::open(&db_path).unwrap(); - let _ = RocksDBBackend::open_with_stale_lock_recovery(&db_path); + #[test] + fn test_waiting_open_succeeds_once_the_holder_releases() { + let temp_dir = TempDir::new().unwrap(); + let db_path = temp_dir.path().to_path_buf(); + let mut holder = RocksDBBackend::open(&db_path).unwrap(); + holder.put(b"k", b"v").unwrap(); + let release = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(200)); + drop(holder); + }); + + let backend = + RocksDBBackend::open_waiting_for_lock(&db_path, Duration::from_secs(10)).unwrap(); + assert_eq!(backend.get(b"k").unwrap(), Some(b"v".to_vec())); + release.join().unwrap(); } #[test] From fee1536be42d1e2412520716a32711859a93ceb0 Mon Sep 17 00:00:00 2001 From: Andrey Vasilevsky Date: Sat, 3 Oct 2026 22:14:06 -0700 Subject: [PATCH 2/2] no-mistakes(review): Disarm load sentinel during lock waits; share cross-project handle --- crates/codegraph-server/src/domain/impact.rs | 38 ++++----- crates/codegraph-server/src/mcp/server.rs | 62 ++++++++------ crates/codegraph-server/src/memory.rs | 17 +++- .../tests/graph_db_recovery_test.rs | 1 - .../codegraph/src/storage/rocksdb_backend.rs | 82 ++++++++++++++++++- 5 files changed, 150 insertions(+), 50 deletions(-) diff --git a/crates/codegraph-server/src/domain/impact.rs b/crates/codegraph-server/src/domain/impact.rs index 2086b97..21024b8 100644 --- a/crates/codegraph-server/src/domain/impact.rs +++ b/crates/codegraph-server/src/domain/impact.rs @@ -336,24 +336,22 @@ fn find_cross_project_consumers( } }; - // Open RocksDB, scan registry, then DROP the connection before per-project - // loading — RocksDB uses exclusive locks, so only one connection at a time. - let entries = { - let rocks = match crate::memory::open_shared_graph_db(&db_path) { - Ok(r) => r, - Err(e) => { - tracing::warn!("[cross-project] Failed to open graph.db: {}", e); - return Vec::new(); - } - }; - match StorageBackend::scan_prefix(&rocks, b"_registry:") { - Ok(e) => e, - Err(e) => { - tracing::warn!("[cross-project] Failed to scan registry: {}", e); - return Vec::new(); - } + // Open RocksDB once for the registry scan and every per-project load, so + // a busy DB costs this query at most one lock wait. The handle drops, and + // the lock is released, when this function returns. + let rocks = match crate::memory::open_shared_graph_db(&db_path) { + Ok(r) => r, + Err(e) => { + tracing::warn!("[cross-project] Failed to open graph.db: {}", e); + return Vec::new(); + } + }; + let entries = match StorageBackend::scan_prefix(&rocks, b"_registry:") { + Ok(e) => e, + Err(e) => { + tracing::warn!("[cross-project] Failed to scan registry: {}", e); + return Vec::new(); } - // rocks dropped here — lock released }; tracing::debug!( @@ -399,11 +397,7 @@ fn find_cross_project_consumers( }) .unwrap_or_else(|| slug.clone()); - let other_rocks = match crate::memory::open_shared_graph_db(&db_path) { - Ok(r) => r, - Err(_) => continue, - }; - let namespaced = NamespacedBackend::new(Box::new(other_rocks), &slug); + let namespaced = NamespacedBackend::new(Box::new(rocks.clone()), &slug); let mut other_graph = match CodeGraph::with_backend(Box::new(namespaced)) { Ok(g) => g, Err(_) => continue, diff --git a/crates/codegraph-server/src/mcp/server.rs b/crates/codegraph-server/src/mcp/server.rs index 6badedc..774b067 100644 --- a/crates/codegraph-server/src/mcp/server.rs +++ b/crates/codegraph-server/src/mcp/server.rs @@ -327,28 +327,15 @@ impl McpBackend { } } - // Mark WHERE we are (telemetry) and arm this process's sentinel for - // the load. The phase guard resets to `serving` on return; the + // Mark WHERE we are (telemetry); the load arms this process's + // sentinel. The phase guard resets to `serving` on return; the // sentinel only clears on a completed load (success or graceful - // error), not on a native AV. The sentinel body carries this - // process's start time so a recycled PID can't impersonate a live - // loader (Windows reuses PIDs aggressively during crash-restart - // churn). + // error) or a lock wait, not on a native AV. The sentinel body + // carries this process's start time so a recycled PID can't + // impersonate a live loader (Windows reuses PIDs aggressively during + // crash-restart churn). let _phase = crate::crash_phase::enter("graph_load"); - let sentinel = db_path - .parent() - .map(|p| p.join(format!("graph.loading.{}", std::process::id()))); - if let Some(s) = &sentinel { - let body = Self::own_start_time() - .map(|t| t.to_string()) - .unwrap_or_default(); - let _ = std::fs::write(s, body); - } - let result = Self::load_persistent_graph_inner(&db_path, slug); - if let Some(s) = &sentinel { - let _ = std::fs::remove_file(s); - } match result { Ok(graph) => Ok(graph), @@ -521,18 +508,43 @@ impl McpBackend { /// detach storage to release the lock. The fallible core of /// [`open_persistent_graph`], split out so the poison-pill wrapper can /// retry it on a fresh DB after quarantining a corrupt one. + /// + /// This process's `graph.loading.` sentinel is armed for every open + /// attempt (open itself can AV on a torn DB) and kept from a successful + /// open through the load, but removed while waiting out another holder: + /// a session killed mid-wait must not leave poison evidence behind. fn load_persistent_graph_inner( db_path: &std::path::Path, slug: &str, ) -> Result { - let rocks = memory::open_shared_graph_db(db_path)?; - let namespaced = NamespacedBackend::new(Box::new(rocks), slug); - let mut graph = CodeGraph::with_backend(Box::new(namespaced))?; + let sentinel = db_path + .parent() + .map(|p| p.join(format!("graph.loading.{}", std::process::id()))); + let body = Self::own_start_time() + .map(|t| t.to_string()) + .unwrap_or_default(); + let arm = || { + if let Some(s) = &sentinel { + let _ = std::fs::write(s, &body); + } + }; + let disarm = || { + if let Some(s) = &sentinel { + let _ = std::fs::remove_file(s); + } + }; - // Detach to release the RocksDB lock — all data is now in memory - graph.detach_storage()?; + let result = memory::open_shared_graph_db_with(db_path, arm, disarm).and_then(|rocks| { + let namespaced = NamespacedBackend::new(Box::new(rocks), slug); + let mut graph = CodeGraph::with_backend(Box::new(namespaced))?; - Ok(graph) + // Detach to release the RocksDB lock — all data is now in memory + graph.detach_storage()?; + + Ok(graph) + }); + disarm(); + result } /// Move a corrupt `graph.db` aside so the next open starts clean. RocksDB diff --git a/crates/codegraph-server/src/memory.rs b/crates/codegraph-server/src/memory.rs index e9ae983..7c3cca6 100644 --- a/crates/codegraph-server/src/memory.rs +++ b/crates/codegraph-server/src/memory.rs @@ -115,7 +115,22 @@ pub(crate) const SHARED_GRAPH_DB_LOCK_WAIT: std::time::Duration = /// handle at a time, so an open that gives up at the first refusal fails /// whenever a sibling session happens to be persisting. pub(crate) fn open_shared_graph_db(path: &Path) -> codegraph::Result { - codegraph::RocksDBBackend::open_waiting_for_lock(path, SHARED_GRAPH_DB_LOCK_WAIT) + open_shared_graph_db_with(path, || {}, || {}) +} + +/// [`open_shared_graph_db`] with the per-attempt hooks of +/// [`codegraph::RocksDBBackend::open_waiting_for_lock_with`]. +pub(crate) fn open_shared_graph_db_with( + path: &Path, + before_attempt: impl FnMut(), + on_locked: impl FnMut(), +) -> codegraph::Result { + codegraph::RocksDBBackend::open_waiting_for_lock_with( + path, + SHARED_GRAPH_DB_LOCK_WAIT, + before_attempt, + on_locked, + ) } /// Redirect the shared graph DB to a brand-new directory by bumping the diff --git a/crates/codegraph-server/tests/graph_db_recovery_test.rs b/crates/codegraph-server/tests/graph_db_recovery_test.rs index 6cf4336..394c94d 100644 --- a/crates/codegraph-server/tests/graph_db_recovery_test.rs +++ b/crates/codegraph-server/tests/graph_db_recovery_test.rs @@ -88,7 +88,6 @@ impl Env { "tools/call", json!({"name": "codegraph_symbol_search", "arguments": {"query": MARKER}}), ); - drop(request); drop(stdin); assert!(child.wait().unwrap().success()); reply.to_string().contains(MARKER) diff --git a/crates/codegraph/src/storage/rocksdb_backend.rs b/crates/codegraph/src/storage/rocksdb_backend.rs index fdbd64f..42b8e90 100644 --- a/crates/codegraph/src/storage/rocksdb_backend.rs +++ b/crates/codegraph/src/storage/rocksdb_backend.rs @@ -72,17 +72,43 @@ impl RocksDBBackend { /// Returns [`GraphError::Locked`] if the database is still held after /// `max_wait`, or [`GraphError::Storage`] for any other open failure. pub fn open_waiting_for_lock>(path: P, max_wait: Duration) -> Result { + Self::open_waiting_for_lock_with(path, max_wait, || {}, || {}) + } + + /// [`Self::open_waiting_for_lock`], calling `before_attempt` right before + /// every open attempt and `on_locked` right after every attempt refused + /// with [`GraphError::Locked`], including the last one. Nothing runs + /// between `on_locked` and the next `before_attempt` except the backoff + /// sleep, so state set up in `before_attempt` and torn down in + /// `on_locked` never exists while waiting. + /// + /// # Errors + /// + /// As [`Self::open_waiting_for_lock`]. + pub fn open_waiting_for_lock_with>( + path: P, + max_wait: Duration, + mut before_attempt: impl FnMut(), + mut on_locked: impl FnMut(), + ) -> Result { let deadline = Instant::now() + max_wait; let mut delay = Duration::from_millis(20); loop { + before_attempt(); match Self::open(path.as_ref()) { Err(GraphError::Locked { .. }) if Instant::now() < deadline => { + on_locked(); std::thread::sleep( delay.min(deadline.saturating_duration_since(Instant::now())), ); delay = (delay * 2).min(Duration::from_millis(500)); } - result => return result, + result => { + if matches!(result, Err(GraphError::Locked { .. })) { + on_locked(); + } + return result; + } } } } @@ -474,6 +500,60 @@ mod tests { assert_eq!(holder.get(b"k").unwrap(), Some(b"v".to_vec())); } + #[test] + fn test_waiting_open_hooks_bracket_every_refused_attempt() { + let temp_dir = TempDir::new().unwrap(); + let db_path = temp_dir.path().to_path_buf(); + let _holder = RocksDBBackend::open(&db_path).unwrap(); + + let attempts = std::cell::Cell::new(0u32); + let refused = std::cell::Cell::new(0u32); + let err = RocksDBBackend::open_waiting_for_lock_with( + &db_path, + Duration::from_millis(300), + || { + assert_eq!(attempts.get(), refused.get(), "attempt began while armed"); + attempts.set(attempts.get() + 1); + }, + || refused.set(refused.get() + 1), + ) + .err() + .unwrap(); + + assert!(matches!(err, GraphError::Locked { .. }), "got {err}"); + assert!( + attempts.get() > 1, + "expected retries, got {}", + attempts.get() + ); + assert_eq!(refused.get(), attempts.get()); + } + + #[test] + fn test_waiting_open_hooks_leave_a_successful_attempt_armed() { + let temp_dir = TempDir::new().unwrap(); + let db_path = temp_dir.path().to_path_buf(); + let holder = RocksDBBackend::open(&db_path).unwrap(); + let release = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(200)); + drop(holder); + }); + + let attempts = std::cell::Cell::new(0u32); + let refused = std::cell::Cell::new(0u32); + RocksDBBackend::open_waiting_for_lock_with( + &db_path, + Duration::from_secs(10), + || attempts.set(attempts.get() + 1), + || refused.set(refused.get() + 1), + ) + .unwrap(); + release.join().unwrap(); + + assert!(refused.get() > 0); + assert_eq!(attempts.get(), refused.get() + 1); + } + #[test] fn test_waiting_open_succeeds_once_the_holder_releases() { let temp_dir = TempDir::new().unwrap();