Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions src/memory/conversations/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,11 +43,11 @@ pub use bus::{
ConversationEventBus, ConversationPersistenceSubscriber,
};
pub use store::{
append_message, delete_thread, ensure_thread, get_messages, list_threads, purge_threads,
update_message, update_thread_labels, update_thread_title, ConversationPurgeStats,
ConversationStore,
append_message, delete_messages_from, delete_thread, ensure_thread, get_messages, list_threads,
purge_threads, update_message, update_thread_labels, update_thread_title,
ConversationPurgeStats, ConversationStore,
};
pub use types::{
ConversationMessage, ConversationMessagePatch, ConversationThread, CreateConversationThread,
CrossThreadHit,
is_deterministic_message_id, reply_run_id, run_reply_message_id, ConversationMessage,
ConversationMessagePatch, ConversationThread, CreateConversationThread, CrossThreadHit,
};
144 changes: 104 additions & 40 deletions src/memory/conversations/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,24 +4,23 @@
//! `threads/<hex(thread_id)>.jsonl` so arbitrary provider ids remain
//! filesystem-safe.
//!
//! All on-disk mutations serialise through a single process-wide mutex so
//! concurrent RPC handlers don't interleave writes.
//! On-disk mutations synchronize at the narrowest safe scope: lifecycle per
//! conversation root, shared metadata per root, and messages per thread.
//!
//! Ported from OpenHuman's `memory_conversations::store`. The behaviour is
//! preserved; the only mechanical changes are dependency substitutions that
//! keep this crate's `Cargo.toml` untouched:
//! Dependency substitutions that keep this crate's `Cargo.toml` untouched:
//! [`std::sync::LazyLock`] for the statics, the local [`hex_encode`] for
//! per-thread filenames, and the hand-rolled temp-write in [`rewrite_jsonl`].
//! They produce the bytes and paths every existing transcript already lives
//! at, so they must not change.
//!
//! - `once_cell::sync::Lazy` → `std::sync::LazyLock` for the process-wide
//! statics.
//! - `hex::encode` → the local [`hex_encode`] helper for per-thread filenames.
//! - `tempfile::NamedTempFile` (a dev-only dependency here) → a write-to-temp +
//! atomic-rename in [`rewrite_jsonl`].
//! - OpenHuman's `log`/`tracing` diagnostics are dropped (this crate has no
//! logging facade wired up).
//! The lock registry ([`locks`]) and deterministic-id idempotency were
//! upstreamed from OpenHuman's host copy of this store.
//!
//! To respect the repo's 500-line-per-file limit the `impl ConversationStore`
//! is split across two child modules — [`ops`] (the public CRUD + search API)
//! and [`index`] (private thread-folding and inverted-index helpers). Both are
//! # File split
//!
//! To respect the repo's file-size limit the `impl ConversationStore` is split
//! across two child modules — [`ops`] (the public CRUD + search API) and
//! [`index`] (private thread-folding and inverted-index helpers). Both are
//! descendant modules of `store`, so they share access to the private statics,
//! constants, log-entry enum, and JSONL helpers defined here.

Expand All @@ -44,6 +43,8 @@ mod ops;

#[path = "store_index.rs"]
mod index;
#[path = "store_locks.rs"]
mod locks;

/// Filename of the append-only thread metadata log, relative to the
/// `memory/conversations` root.
Expand All @@ -52,10 +53,6 @@ pub(super) const THREADS_FILENAME: &str = "threads.jsonl";
/// per-thread message JSONL files, named `<hex(thread_id)>.jsonl`.
pub(super) const THREAD_MESSAGES_DIR: &str = "threads";

/// Serialises every on-disk mutation so concurrent handlers can't interleave
/// writes to `threads.jsonl` or the per-thread message logs.
static CONVERSATION_STORE_LOCK: LazyLock<Mutex<()>> = LazyLock::new(|| Mutex::new(()));

/// Per-workspace inverted index cache. Keyed by the workspace's
/// `memory/conversations` root so multiple `ConversationStore` clones
/// pointing at the same workspace share one index. The cache outlives
Expand All @@ -67,28 +64,26 @@ static CONVERSATION_STORE_LOCK: LazyLock<Mutex<()>> = LazyLock::new(|| Mutex::ne
///
/// # Lock ordering
///
/// When BOTH `CONVERSATION_STORE_LOCK` and `CONVERSATION_INDEX_CACHE`
/// must be held simultaneously, `CONVERSATION_STORE_LOCK` MUST be
/// acquired first. This applies to `append_message` (writes JSONL then
/// updates the warm index) and `with_index` (caller holds the outer
/// lock, then takes the cache lock to run the search closure).
/// Every operation first takes the root lifecycle lock. Message operations
/// then take their per-thread lock and briefly take the metadata lock when
/// they must inspect or append `threads.jsonl`. When metadata and the index
/// cache are both needed, metadata is acquired first. No code may acquire a
/// thread or metadata lock while holding `CONVERSATION_INDEX_CACHE`.
///
/// `prime_index_if_cold` minimises shared locking. It may hold both
/// locks only momentarily, and always in the `CONVERSATION_STORE_LOCK`
/// → `CONVERSATION_INDEX_CACHE` order above: while holding the outer
/// lock to snapshot live thread IDs via `thread_index_unlocked`
/// metadata and index locks only momentarily, and always in the metadata
/// → `CONVERSATION_INDEX_CACHE` order above: while holding metadata
/// to snapshot live thread IDs via `thread_index_unlocked`
/// (header-only, no per-thread I/O) it re-checks the cache once. It then
/// releases `CONVERSATION_STORE_LOCK` before reading per-thread JSONL
/// content (no lock held) and finally acquires `CONVERSATION_INDEX_CACHE`
/// alone to insert the built index. It never holds both across the slow
/// JSONL walk, and neither operation calls back into a function that
/// would acquire the other lock.
/// releases metadata before reading each transcript under that thread's own
/// lock and finally acquires `CONVERSATION_INDEX_CACHE` alone to insert the
/// built index. It never holds both across the slow JSONL walk.
///
/// `list_threads_unlocked` MUST NOT be used inside the locked snapshot —
/// it calls `measure_messages_unlocked` per legacy thread (no Stats
/// history), which reads every per-thread JSONL file and appends a
/// `Stats` entry to `threads.jsonl`, reintroducing the multi-second
/// stall under the outer lock that this design was built to avoid.
/// stall under the shared metadata lock that this design was built to avoid.
static CONVERSATION_INDEX_CACHE: LazyLock<Mutex<HashMap<PathBuf, InvertedIndex>>> =
LazyLock::new(|| Mutex::new(HashMap::new()));

Expand All @@ -104,20 +99,39 @@ pub struct ConversationPurgeStats {
/// Workspace-rooted handle that reads and writes the JSONL conversation log.
#[derive(Debug, Clone)]
pub struct ConversationStore {
workspace_dir: PathBuf,
root_dir: PathBuf,
locks: std::sync::Arc<locks::StoreLocks>,
}

impl ConversationStore {
/// Construct a store rooted at the given workspace directory.
///
/// The conversation root is derived from `workspace_dir` alone
/// (`<workspace>/memory/conversations`, see `root_dir` in `store_index.rs`).
pub fn new(workspace_dir: PathBuf) -> Self {
Self { workspace_dir }
let root = locks::normalized_root(&workspace_dir.join("memory").join("conversations"));
let locks = locks::for_root(&root);
Self {
root_dir: root,
locks,
}
}

/// Construct a store rooted at the engine's configured workspace
/// ([`MemoryConfig::workspace`](crate::memory::config::MemoryConfig)).
pub fn from_config(config: &crate::memory::config::MemoryConfig) -> Self {
Self::new(config.workspace.clone())
}

#[cfg(test)]
pub(super) fn lock_identity_for_test(&self) -> usize {
std::sync::Arc::as_ptr(&self.locks) as usize
}

#[cfg(test)]
pub(super) fn thread_lock_count_for_test(&self) -> usize {
self.locks.thread_count()
}
}

/// One line in `threads.jsonl`. The append-only log is folded into the current
Expand Down Expand Up @@ -220,8 +234,13 @@ pub(super) fn normalize_labels(labels: Vec<String>) -> Vec<String> {
}

/// Lowercase hex-encode bytes — used to derive a filesystem-safe per-thread
/// messages filename from an arbitrary thread id. Replaces OpenHuman's use of
/// the `hex` crate, which is not a dependency of this crate.
/// messages filename from an arbitrary thread id.
///
/// This exists because the memory engine had no `hex` dependency; this crate
/// does, so `hex::encode` would work here. It is kept anyway: this is the
/// function that decides which file a user's transcript is read from and
/// written to, and swapping it is only provably safe, never *obviously* safe.
/// The three lines are cheaper than the argument.
pub(super) fn hex_encode(bytes: &[u8]) -> String {
const HEX: [u8; 16] = *b"0123456789abcdef";
let mut out = String::with_capacity(bytes.len() * 2);
Expand Down Expand Up @@ -258,6 +277,36 @@ where
Ok(items)
}

/// Find one message in a thread's JSONL log by id, without materializing the
/// whole transcript.
///
/// Only the lines whose raw text carries the quoted id are deserialized, so a
/// lookup costs one parse rather than one per stored message; a line that
/// merely quotes the id inside its own content is rejected by the `id` check.
/// Mirrors [`read_jsonl`]'s tolerance of blank and corrupt lines.
pub(super) fn find_message_by_id(
path: &Path,
id: &str,
) -> Result<Option<ConversationMessage>, String> {
if !path.exists() {
return Ok(None);
}
let needle = serde_json::to_string(id).map_err(|e| format!("encode message id {id}: {e}"))?;
let file = File::open(path).map_err(|e| format!("open {}: {e}", path.display()))?;
for (line_no, line) in BufReader::new(file).lines().enumerate() {
let line =
line.map_err(|e| format!("read {} line {}: {e}", path.display(), line_no + 1))?;
if !line.contains(&needle) {
continue;
}
match serde_json::from_str::<ConversationMessage>(&line) {
Ok(message) if message.id == id => return Ok(Some(message)),
_ => continue,
}
}
Ok(None)
}

/// Append one serialized value as a JSONL line, fsync'd before returning.
pub(super) fn append_jsonl<T>(path: &Path, value: &T) -> Result<(), String>
where
Expand All @@ -284,9 +333,15 @@ where
/// Atomically rewrite `path` with `values`, one JSON object per line.
///
/// Writes to a sibling temp file then renames over the target so a crash
/// mid-write never leaves a partially-written transcript. Replaces
/// OpenHuman's `tempfile::NamedTempFile`, which is a dev-only dependency in
/// this crate.
/// mid-write never leaves a partially-written transcript.
///
/// This hand-rolls what OpenHuman originally got from
/// `tempfile::NamedTempFile`, because the memory engine carried `tempfile` as
/// a dev-dependency only. This crate has it in full, so the substitution is no
/// longer forced — but reverting it would change the temp file's name, its
/// permissions, and which side deletes it when a write fails. That is a change
/// to the crash-safety path for a user's transcript, and it belongs in a
/// change that is about that, not in a module move.
pub(super) fn rewrite_jsonl<T>(path: &Path, values: &[T]) -> Result<(), String>
where
T: serde::Serialize,
Expand Down Expand Up @@ -381,6 +436,15 @@ pub fn update_message(
ConversationStore::new(workspace_dir).update_message(thread_id, message_id, patch)
}

/// Free-function shim around [`ConversationStore::delete_messages_from`].
pub fn delete_messages_from(
workspace_dir: PathBuf,
thread_id: &str,
message_id: &str,
) -> Result<Option<usize>, String> {
ConversationStore::new(workspace_dir).delete_messages_from(thread_id, message_id)
}

/// Free-function shim around [`ConversationStore::purge_threads`].
pub fn purge_threads(workspace_dir: PathBuf) -> Result<ConversationPurgeStats, String> {
ConversationStore::new(workspace_dir).purge_threads()
Expand Down
104 changes: 104 additions & 0 deletions src/memory/conversations/store_concurrency_tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
use super::*;

#[test]
fn one_hundred_agent_threads_append_without_loss_or_corruption() {
use std::sync::{Arc, Barrier};

let temp = TempDir::new().unwrap();
let store = ConversationStore::new(temp.path().to_path_buf());
let created_at = "2026-09-10T00:00:00Z".to_string();

for index in 0..100 {
store
.ensure_thread(CreateConversationThread {
parent_thread_id: None,
id: format!("agent-{index}"),
title: format!("Agent {index}"),
created_at: created_at.clone(),
labels: None,
personality_id: None,
})
.unwrap();
}

let barrier = Arc::new(Barrier::new(101));
let mut writers = Vec::with_capacity(100);
for index in 0..100 {
let store = store.clone();
let barrier = Arc::clone(&barrier);
let created_at = created_at.clone();
writers.push(std::thread::spawn(move || {
barrier.wait();
store.append_message(
&format!("agent-{index}"),
ConversationMessage {
id: format!("message-{index}"),
content: format!("reply from agent {index}"),
message_type: "text".to_string(),
extra_metadata: serde_json::json!({}),
sender: "assistant".to_string(),
created_at,
},
)
}));
}
barrier.wait();

for writer in writers {
writer.join().expect("writer panicked").expect("append");
}

let threads = store.list_threads().unwrap();
assert_eq!(threads.len(), 100);
assert!(threads.iter().all(|thread| thread.message_count == 1));
for index in 0..100 {
let messages = store.get_messages(&format!("agent-{index}")).unwrap();
assert_eq!(messages.len(), 1);
assert_eq!(messages[0].id, format!("message-{index}"));
}
}

#[test]
fn one_unreadable_transcript_does_not_stop_other_repairs() {
let temp = TempDir::new().unwrap();
let store = ConversationStore::new(temp.path().to_path_buf());
let created_at = "2026-09-10T00:00:00Z".to_string();

for id in ["a-unreadable", "z-readable"] {
store
.ensure_thread(CreateConversationThread {
parent_thread_id: None,
id: id.to_string(),
title: id.to_string(),
created_at: created_at.clone(),
labels: None,
personality_id: None,
})
.unwrap();
}

std::fs::create_dir_all(store.thread_messages_path("a-unreadable")).unwrap();
append_jsonl(
&store.thread_messages_path("z-readable"),
&ConversationMessage {
id: "message".to_string(),
content: "recover me".to_string(),
message_type: "text".to_string(),
extra_metadata: serde_json::json!({}),
sender: "assistant".to_string(),
created_at,
},
)
.unwrap();

let threads = store.list_threads().unwrap();
let readable = threads
.iter()
.find(|thread| thread.id == "z-readable")
.unwrap();
assert_eq!(readable.message_count, 1);

let folded = store.thread_index_unlocked().unwrap();
assert_eq!(folded["z-readable"].message_count, Some(1));
assert_eq!(folded["a-unreadable"].message_count, None);
}
Loading
Loading