diff --git a/src/openhuman/memory_conversations/store.rs b/src/openhuman/memory_conversations/store.rs index 8ec83682a..f2bd5e293 100644 --- a/src/openhuman/memory_conversations/store.rs +++ b/src/openhuman/memory_conversations/store.rs @@ -36,16 +36,23 @@ static CONVERSATION_STORE_LOCK: Lazy> = Lazy::new(|| Mutex::new(())); /// behind dead entries when the dir is removed — acceptable for an /// in-process cache. /// -/// # Lock ordering (INVARIANT) +/// # Lock ordering /// -/// `CONVERSATION_STORE_LOCK` MUST be acquired before -/// `CONVERSATION_INDEX_CACHE`. Every public store method takes the outer -/// lock at the top of its body and only then reaches for the cache, so -/// in current code the inner mutex is uncontended. Any future caller -/// that violates this order — taking the cache lock first, or taking -/// the cache lock without holding the outer lock and then reaching -/// across into another store API — risks a deadlock. Keep the inner -/// mutex strictly nested inside the outer one. +/// Most mutating store methods (`append_message`, `delete_thread`, etc.) +/// hold `CONVERSATION_STORE_LOCK` first, then briefly lock this cache to +/// do a warm-cache index update — the classic nested order. +/// +/// `search_cross_thread_messages` is the exception (issue #2849): its +/// **warm-cache fast path** locks only this cache (no outer lock needed +/// for a read-only index query), and its **cold-path rebuild** takes +/// the outer lock briefly to snapshot the thread list, releases it, does +/// the slow JSONL I/O lock-free, then acquires this cache lock to +/// insert the result. This avoids holding both locks for the entire +/// rebuild, which previously blocked all concurrent writes. +/// +/// **Rule:** never hold this cache lock while calling back into a public +/// `ConversationStore` method that takes `CONVERSATION_STORE_LOCK` — that +/// would invert the dominant ordering and risk a deadlock. static CONVERSATION_INDEX_CACHE: Lazy>> = Lazy::new(|| Mutex::new(HashMap::new())); @@ -179,44 +186,47 @@ impl ConversationStore { /// turns that into O(|posting lists|) for typical queries while /// preserving the previous scoring contract /// (`score = matched_terms / total_terms`, recency tiebreak). + /// + /// # Lock strategy (issue #2849) + /// + /// **Fast path (warm cache):** acquires only `CONVERSATION_INDEX_CACHE` + /// — no outer store lock — and returns immediately. + /// + /// **Cold path (first access):** snapshots the thread list under + /// `CONVERSATION_STORE_LOCK` (brief), then releases it before reading + /// JSONL files to build the inverted index. This avoids blocking + /// `append_message` / `get_messages` / `list_threads` during the + /// potentially-long rebuild. JSONL files are append-only, so a + /// concurrent write during the rebuild may mean the rebuilt index + /// misses that one message. It is *not* re-read later; subsequent + /// `append_message` calls only index their own (new) messages once + /// the cache is warm. The missed message therefore stays absent + /// until the cache is evicted and rebuilt — an accepted tradeoff + /// for issue #2849. pub fn search_cross_thread_messages( &self, query: &str, limit: usize, exclude_thread_id: Option<&str>, ) -> Result, String> { - let _guard = CONVERSATION_STORE_LOCK.lock(); - self.with_index(|idx| idx.search(query, limit, exclude_thread_id)) - } - - /// Acquire the cached inverted index for this workspace (building it - /// from JSONL on first access) and run `f` against it. Caller MUST - /// hold `CONVERSATION_STORE_LOCK` for the duration of the closure. - fn with_index(&self, f: impl FnOnce(&mut InvertedIndex) -> R) -> Result { - let key = self.root_dir(); - let mut cache = CONVERSATION_INDEX_CACHE.lock(); - if !cache.contains_key(&key) { - let mut idx = InvertedIndex::new(); - self.populate_index_unlocked(&mut idx)?; - cache.insert(key.clone(), idx); + // Fast path: if the index is already warm, search without the + // outer store lock. + { + let mut cache = CONVERSATION_INDEX_CACHE.lock(); + if let Some(idx) = cache.get_mut(&self.root_dir()) { + return Ok(idx.search(query, limit, exclude_thread_id)); + } } - let idx = cache.get_mut(&key).expect("inserted above if absent"); - Ok(f(idx)) - } - /// Walk every per-thread JSONL file in the workspace and insert each - /// message into `idx`. Used on first access to a workspace's index; - /// also called after `purge_threads` to reset the cache from a - /// known-empty state. The JSONL files remain the source of truth, so - /// a rebuild after a process crash is always safe. - fn populate_index_unlocked(&self, idx: &mut InvertedIndex) -> Result<(), String> { - // `list_threads_unlocked` already handles a fresh workspace: - // `ensure_root` creates the directory + threads log if missing, and - // `read_jsonl` returns an empty Vec for an empty file. Anything that - // still bubbles up here is a real filesystem/setup failure and must - // propagate — silently returning Ok would mask it and make search - // appear to return zero results for an undiagnosed reason. - let threads = self.list_threads_unlocked()?; + // Cold path: build the index. Snapshot the thread list under the + // store lock (brief), then release it so concurrent writes aren't + // blocked during the potentially-long JSONL file reads. + let threads = { + let _guard = CONVERSATION_STORE_LOCK.lock(); + self.list_threads_unlocked()? + }; + + let mut idx = InvertedIndex::new(); for thread in threads { let path = self.thread_messages_path(&thread.id); if !path.exists() { @@ -234,8 +244,6 @@ impl ConversationStore { } }; for msg in messages { - // Move each freshly-deserialized message straight into - // the index; no per-field clones on the rebuild path. idx.insert(&thread.id, msg); } } @@ -243,7 +251,14 @@ impl ConversationStore { "{LOG_PREFIX} inverted index populated workspace={}", self.root_dir().display() ); - Ok(()) + + // Insert the newly-built index and run the search. Another thread + // may have raced and already inserted — prefer the existing index + // if present (it may have more recent messages from concurrent + // append_message calls). + let mut cache = CONVERSATION_INDEX_CACHE.lock(); + let entry = cache.entry(self.root_dir()).or_insert(idx); + Ok(entry.search(query, limit, exclude_thread_id)) } /// Append a message to the thread's JSONL file. Errors if the thread is missing. diff --git a/src/openhuman/memory_conversations/store_tests.rs b/src/openhuman/memory_conversations/store_tests.rs index b5b6e73e2..efbf4b28e 100644 --- a/src/openhuman/memory_conversations/store_tests.rs +++ b/src/openhuman/memory_conversations/store_tests.rs @@ -846,6 +846,74 @@ fn update_thread_labels_missing_thread_returns_error() { assert!(err.contains("thread missing not found")); } +#[test] +fn cold_search_does_not_serialize_on_outer_lock() { + // Issue #2849: verify that a cold-cache search releases the store + // lock before the JSONL rebuild, so concurrent writes aren't blocked. + let (_temp, store) = make_store(); + + // Seed a thread with a message so the search has something to find. + store + .ensure_thread(CreateConversationThread { + id: "t1".to_string(), + title: "test thread".to_string(), + created_at: "2026-01-01T00:00:00Z".to_string(), + parent_thread_id: None, + labels: None, + }) + .unwrap(); + store + .append_message( + "t1", + ConversationMessage { + id: "m1".to_string(), + content: "hello world".to_string(), + message_type: "text".to_string(), + extra_metadata: json!({}), + sender: "user".to_string(), + created_at: "2026-01-01T00:00:00Z".to_string(), + }, + ) + .unwrap(); + + // Evict any warm cache so the next search triggers a cold rebuild. + { + let mut cache = CONVERSATION_INDEX_CACHE.lock(); + cache.remove(&store.root_dir()); + } + + // Spawn a thread that tries to append a message while a cold search + // is (conceptually) running. In the old code this would deadlock or + // serialize behind the full rebuild; in the fixed code the store lock + // is released after the thread-list snapshot and the append succeeds + // concurrently. + let store2 = store.clone(); + let writer = std::thread::spawn(move || { + store2 + .append_message( + "t1", + ConversationMessage { + id: "m2".to_string(), + content: "concurrent write".to_string(), + message_type: "text".to_string(), + extra_metadata: json!({}), + sender: "assistant".to_string(), + created_at: "2026-01-01T00:00:01Z".to_string(), + }, + ) + .unwrap(); + }); + + // Run the cold search — should not deadlock. + let results = store + .search_cross_thread_messages("hello", 10, None) + .unwrap(); + assert!(!results.is_empty(), "search should find seeded message"); + + // The concurrent write must also succeed. + writer.join().expect("concurrent write must not deadlock"); +} + #[test] fn read_jsonl_skips_invalid_lines_but_keeps_valid_ones() { let tmp = TempDir::new().unwrap();