use std::collections::HashMap;
use std::fs::{File, OpenOptions};
use std::io::{BufRead, BufReader, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Mutex, PoisonError};
use anyhow::{Context, Result};
use mermaid_domain::{
ConversationHistory, SESSION_EVENT_FORMAT_VERSION, SessionEvent, SessionEventLine,
SessionScalars, fold_session,
};
use mermaid_model::models::{ChatMessage, MessageRole};
const MAX_LOG_BYTES: u64 = 128 * 1024 * 1024;
static CONFLICT_COUNTER: AtomicU64 = AtomicU64::new(0);
const SCREENSHOT_ELIDED_MARKER: &str = "\n[screenshot not persisted]";
pub struct EventLog {
dir: PathBuf,
next_seq: Mutex<HashMap<String, Cursor>>,
diverted: Mutex<HashMap<String, PathBuf>>,
}
#[derive(Debug, Clone, Copy)]
struct Cursor {
next_seq: u64,
len: u64,
}
impl EventLog {
pub(crate) fn new(dir: PathBuf) -> Self {
Self {
dir,
next_seq: Mutex::new(HashMap::new()),
diverted: Mutex::new(HashMap::new()),
}
}
pub(crate) fn path_for(&self, id: &str) -> PathBuf {
self.dir.join(format!("{id}.jsonl"))
}
pub(crate) fn append(
&self,
snapshot: &ConversationHistory,
events: &[SessionEvent],
) -> Result<()> {
if snapshot.messages().is_empty() {
return Ok(());
}
let path = self.active_path(&snapshot.id);
let result = self.append_inner(&path, snapshot, events);
if result.is_err() {
self.next_seq
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remove(&snapshot.id);
}
result
}
fn append_inner(
&self,
path: &Path,
snapshot: &ConversationHistory,
events: &[SessionEvent],
) -> Result<()> {
let diverted = !path.exists() || self.diverted_on_conflict(&snapshot.id, path);
let path = &if diverted {
self.active_path(&snapshot.id)
} else {
path.to_path_buf()
};
let creating = !path.exists();
let batch: Vec<SessionEvent> = if creating {
let mut seeded = backfill_events(snapshot);
seeded.extend(events.iter().filter(|e| is_idempotent(e)).cloned());
seeded
} else {
events.to_vec()
};
if batch.is_empty() {
return Ok(());
}
let mut seq = if creating {
0
} else {
self.next_seq_for(&snapshot.id, path)?
};
let mut file = open_append(path)?;
let ts = chrono::Local::now();
for event in &batch {
let Some(event) = sanitize_event(event) else {
continue;
};
let line = SessionEventLine {
v: SESSION_EVENT_FORMAT_VERSION,
seq,
ts,
event,
};
let mut value = serde_json::to_value(&line).context("serialize session event")?;
mermaid_model::utils::redact_json(&mut value);
writeln!(file, "{value}").context("append session event line")?;
seq += 1;
}
file.flush().context("flush session event log")?;
let len = std::fs::metadata(path).map_or(0, |meta| meta.len());
self.next_seq
.lock()
.unwrap_or_else(PoisonError::into_inner)
.insert(snapshot.id.clone(), Cursor { next_seq: seq, len });
Ok(())
}
fn diverted_on_conflict(&self, id: &str, path: &Path) -> bool {
let Some(cursor) = self
.next_seq
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(id)
.copied()
else {
return false;
};
let actual = std::fs::metadata(path).map_or(0, |meta| meta.len());
if actual == cursor.len {
return false;
}
let sibling = self.conflict_path(id);
tracing::warn!(
id,
log = %path.display(),
conflict = %sibling.display(),
expected_len = cursor.len,
actual_len = actual,
"another mermaid appended to this session's log; this process continues in a .conflict sibling"
);
self.diverted
.lock()
.unwrap_or_else(PoisonError::into_inner)
.insert(id.to_string(), sibling);
self.next_seq
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remove(id);
true
}
fn active_path(&self, id: &str) -> PathBuf {
self.diverted
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(id)
.cloned()
.unwrap_or_else(|| self.path_for(id))
}
fn conflict_path(&self, id: &str) -> PathBuf {
let n = CONFLICT_COUNTER.fetch_add(1, Ordering::Relaxed);
self.dir.join(format!(
"{}.{}.{}.conflict.jsonl",
id,
std::process::id(),
n
))
}
fn next_seq_for(&self, id: &str, path: &Path) -> Result<u64> {
if let Some(seq) = self
.next_seq
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(id)
{
return Ok(seq.next_seq);
}
let Ok(meta) = std::fs::metadata(path) else {
return Ok(0);
};
anyhow::ensure!(
meta.is_file(),
"session event log path {} is not a file",
path.display()
);
let file = File::open(path)
.with_context(|| format!("open {} to derive the seq cursor", path.display()))?;
let mut reader = BufReader::new(file);
let mut line = Vec::new();
let mut seq = 0u64;
loop {
line.clear();
let read = reader
.read_until(b'\n', &mut line)
.with_context(|| format!("scan {} for the seq cursor", path.display()))?;
if read == 0 {
return Ok(seq);
}
seq += 1;
}
}
pub(crate) fn exists(&self, id: &str) -> bool {
self.path_for(id).is_file()
}
pub(crate) fn checkpoint_seq(&self, id: &str) -> Option<u64> {
self.next_seq
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(id)
.and_then(|cursor| cursor.next_seq.checked_sub(1))
}
pub(crate) fn read_events(
&self,
id: &str,
after_seq: Option<u64>,
) -> Result<Option<(Vec<SessionEvent>, u64)>> {
let path = self.path_for(id);
let Ok(meta) = std::fs::metadata(&path) else {
return Ok(None);
};
if !meta.is_file() {
return Ok(None);
}
if meta.len() > MAX_LOG_BYTES {
tracing::warn!(path = %path.display(), "session event log over the read cap; not reading");
return Ok(None);
}
let file =
File::open(&path).with_context(|| format!("open {} for fold", path.display()))?;
let mut events = Vec::new();
let mut highest = 0u64;
for line in BufReader::new(file).lines() {
let raw = line.context("read session event line")?;
let parsed: SessionEventLine = match serde_json::from_str(&raw) {
Ok(parsed) => parsed,
Err(error) => {
tracing::warn!(path = %path.display(), %error, "malformed session event line; reading the prefix only");
break;
},
};
if parsed.v > SESSION_EVENT_FORMAT_VERSION {
tracing::warn!(
path = %path.display(),
version = parsed.v,
"session event log written by a newer mermaid; not reading"
);
return Ok(None);
}
highest = highest.max(parsed.seq);
if after_seq.is_some_and(|seq| parsed.seq <= seq) {
continue;
}
events.push(parsed.event);
}
Ok(Some((events, highest)))
}
pub(crate) fn fold(&self, id: &str) -> Result<Option<ConversationHistory>> {
let Some((events, _)) = self.read_events(id, None)? else {
return Ok(None);
};
Ok(fold_session(events))
}
pub(crate) fn replay_onto(
&self,
id: &str,
mut checkpoint: ConversationHistory,
checkpoint_seq: u64,
) -> Result<Option<ConversationHistory>> {
let Some((events, highest)) = self.read_events(id, Some(checkpoint_seq))? else {
return Ok(None);
};
if highest < checkpoint_seq {
tracing::warn!(
id,
checkpoint_seq,
highest,
"checkpoint is ahead of its log; folding from zero instead"
);
return Ok(None);
}
mermaid_domain::replay_events(&mut checkpoint, events);
Ok(Some(checkpoint))
}
}
fn backfill_events(snapshot: &ConversationHistory) -> Vec<SessionEvent> {
let mut events = vec![SessionEvent::Started {
session_id: snapshot.id.clone(),
project_path: snapshot.project_path.clone(),
model_id: snapshot.model_name.clone(),
created_at: snapshot.created_at,
forked_from: snapshot.forked_from.clone(),
parent_session: snapshot.parent_session.clone(),
}];
if !snapshot.messages().is_empty() {
events.push(SessionEvent::Reset {
at: snapshot.updated_at,
messages: snapshot.messages().to_vec(),
});
}
for text in &snapshot.input_history {
events.push(SessionEvent::Input { text: text.clone() });
}
events.push(SessionEvent::State(Box::new(SessionScalars::of(snapshot))));
events.push(SessionEvent::Tasks {
store: snapshot.tasks.clone(),
});
events
}
const fn is_idempotent(event: &SessionEvent) -> bool {
match event {
SessionEvent::Compaction { .. }
| SessionEvent::Reset { .. }
| SessionEvent::State(_)
| SessionEvent::Tasks { .. } => true,
SessionEvent::Started { .. }
| SessionEvent::Message { .. }
| SessionEvent::InsertedBeforeLast { .. }
| SessionEvent::Action { .. }
| SessionEvent::Image { .. }
| SessionEvent::Input { .. } => false,
}
}
fn sanitize_event(event: &SessionEvent) -> Option<SessionEvent> {
match event {
SessionEvent::Image { .. } => None,
SessionEvent::Message { message } => Some(SessionEvent::Message {
message: strip_screenshot(message),
}),
SessionEvent::InsertedBeforeLast { message } => Some(SessionEvent::InsertedBeforeLast {
message: strip_screenshot(message),
}),
SessionEvent::Compaction {
at,
record,
replacement,
} => Some(SessionEvent::Compaction {
at: *at,
record: record.clone(),
replacement: replacement.iter().map(strip_screenshot).collect(),
}),
SessionEvent::Reset { at, messages } => Some(SessionEvent::Reset {
at: *at,
messages: messages.iter().map(strip_screenshot).collect(),
}),
event @ (SessionEvent::Started { .. }
| SessionEvent::Action { .. }
| SessionEvent::State(_)
| SessionEvent::Input { .. }
| SessionEvent::Tasks { .. }) => Some(event.clone()),
}
}
fn strip_screenshot(message: &ChatMessage) -> ChatMessage {
if message.role == MessageRole::User || message.images.is_none() {
return message.clone();
}
let mut stripped = message.clone();
stripped.images = None;
if !stripped.content.ends_with(SCREENSHOT_ELIDED_MARKER) {
stripped.content.push_str(SCREENSHOT_ELIDED_MARKER);
}
stripped
}
fn open_append(path: &Path) -> Result<File> {
let mut opts = OpenOptions::new();
opts.create(true).append(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
opts.mode(0o600);
}
opts.open(path)
.with_context(|| format!("open {} for event append", path.display()))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::session::ConversationManager;
use mermaid_domain::{Config, State};
fn fixed_ts() -> chrono::DateTime<chrono::Local> {
chrono::DateTime::parse_from_rfc3339("2026-07-02T12:00:00.123+00:00")
.unwrap()
.with_timezone(&chrono::Local)
}
fn temp_root(name: &str) -> PathBuf {
let dir =
std::env::temp_dir().join(format!("mermaid_event_log_{name}_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
dir
}
fn driven_state(root: &Path) -> State {
let mut state = State::new(
Config::default(),
root.to_path_buf(),
"ollama/test".to_string(),
fixed_ts(),
std::env::temp_dir(),
);
state
.session
.append(ChatMessage::user("hello event log"), fixed_ts());
state.session.record_input("hello event log".to_string());
state
.session
.append(ChatMessage::assistant("hi there"), fixed_ts());
state
}
fn json_of(c: &ConversationHistory) -> serde_json::Value {
serde_json::to_value(c).expect("serializes")
}
#[test]
fn append_then_fold_equals_the_loaded_snapshot() {
let root = temp_root("roundtrip");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
manager.save_conversation(&snapshot).unwrap();
let loaded = manager.load_conversation(&snapshot.id).unwrap();
let folded = manager
.fold_conversation_from_log(&snapshot.id)
.unwrap()
.expect("log folds");
assert_eq!(
json_of(&folded),
json_of(&loaded),
"fold(disk log) must equal the loaded snapshot"
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn corrupt_snapshot_recovers_from_the_log() {
let root = temp_root("recovery");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
manager.save_conversation(&snapshot).unwrap();
let json_path = manager
.conversations_dir()
.join(format!("{}.json", snapshot.id));
std::fs::write(&json_path, b"{ torn").unwrap();
let recovered = manager
.load_conversation(&snapshot.id)
.expect("recovery must kick in");
assert_eq!(recovered.id, snapshot.id);
assert_eq!(recovered.messages().len(), 2);
assert_eq!(recovered.title, snapshot.title);
std::fs::remove_file(&json_path).unwrap();
let recovered = manager
.load_conversation(&snapshot.id)
.expect("missing snapshot recovers");
assert_eq!(recovered.messages().len(), 2);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn backfill_materializes_once_and_folds_exactly() {
let root = temp_root("backfill");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let _ = state.session.drain_events(&snapshot);
manager.save_conversation(&snapshot).unwrap();
let log_path = manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id));
assert!(!log_path.exists());
manager.append_session_events(&snapshot, &[]).unwrap();
assert!(log_path.exists());
let lines_after_backfill = std::fs::read_to_string(&log_path).unwrap().lines().count();
let folded = manager
.fold_conversation_from_log(&snapshot.id)
.unwrap()
.expect("backfill folds");
assert_eq!(
json_of(&folded),
json_of(&manager.load_conversation(&snapshot.id).unwrap()),
"a compaction-free backfill folds to the snapshot exactly"
);
manager.append_session_events(&snapshot, &[]).unwrap();
assert_eq!(
std::fs::read_to_string(&log_path).unwrap().lines().count(),
lines_after_backfill
);
state
.session
.append(ChatMessage::user("post-upgrade"), fixed_ts());
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
let raw = std::fs::read_to_string(&log_path).unwrap();
let seqs: Vec<u64> = raw
.lines()
.map(|l| serde_json::from_str::<SessionEventLine>(l).unwrap().seq)
.collect();
let expected: Vec<u64> = (0..seqs.len() as u64).collect();
assert_eq!(seqs, expected, "seq is gapless and monotonic");
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn creating_a_log_keeps_idempotent_events_and_drops_additive_ones() {
let root = temp_root("idempotent_carry");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let mut events = state.session.drain_events(&snapshot);
events.push(SessionEvent::Compaction {
at: fixed_ts(),
record: mermaid_domain::CompactionEvent {
id: "c-first".to_string(),
trigger: mermaid_domain::CompactionTrigger::Manual,
created_at: fixed_ts(),
before_tokens: 100,
after_tokens: 10,
archived_message_count: 4,
preserved_message_count: 2,
preserved_turn_count: 1,
summary_tokens: 5,
duration_secs: 0.1,
review_status: mermaid_domain::CompactionReviewStatus::Reviewed,
review_error: None,
focus: None,
archive_path: None,
},
replacement: snapshot.messages().to_vec(),
});
assert!(
events
.iter()
.any(|e| matches!(e, SessionEvent::Message { .. })),
"the fixture must carry additive events for this to prove anything"
);
manager.append_session_events(&snapshot, &events).unwrap();
let folded = manager
.fold_conversation_from_log(&snapshot.id)
.unwrap()
.expect("folds");
assert_eq!(
folded.messages().len(),
snapshot.messages().len(),
"additive events must not double the backfilled transcript"
);
assert_eq!(
folded.compactions.len(),
1,
"the compaction boundary must survive log creation"
);
assert_eq!(folded.compactions[0].id, "c-first");
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn seq_continues_across_manager_instances() {
let root = temp_root("seq_scan");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
let other = ConversationManager::new(&root).unwrap();
state
.session
.append(ChatMessage::user("second process"), fixed_ts());
let snapshot = state.session.snapshot_conversation();
let more = state.session.drain_events(&snapshot);
other.append_session_events(&snapshot, &more).unwrap();
let raw = std::fs::read_to_string(
manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id)),
)
.unwrap();
let seqs: Vec<u64> = raw
.lines()
.map(|l| serde_json::from_str::<SessionEventLine>(l).unwrap().seq)
.collect();
let expected: Vec<u64> = (0..seqs.len() as u64).collect();
assert_eq!(seqs, expected);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn screenshots_and_secrets_never_reach_the_log() {
let root = temp_root("sanitize");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
state.session.attach_image("SHOTBYTES64".to_string());
state.session.append(
ChatMessage::assistant("OPENAI_API_KEY=sk-abcdefghijklmnop1234"),
fixed_ts(),
);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
let raw = std::fs::read_to_string(
manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id)),
)
.unwrap();
assert!(!raw.contains("SHOTBYTES64"), "screenshot leaked: {raw}");
assert!(
!raw.contains("sk-abcdefghijklmnop1234"),
"secret leaked: {raw}"
);
assert!(raw.contains("[REDACTED]"), "redaction marker expected");
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn torn_tail_line_folds_the_prefix() {
let root = temp_root("torn");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
let log_path = manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id));
let mut raw = std::fs::read_to_string(&log_path).unwrap();
raw.push_str("{\"v\":1,\"seq\":99,\"ts\":\"2026-"); std::fs::write(&log_path, raw).unwrap();
let folded = manager
.fold_conversation_from_log(&snapshot.id)
.unwrap()
.expect("prefix folds");
assert_eq!(folded.messages().len(), 2, "the pre-tear state survives");
let _ = std::fs::remove_dir_all(&root);
}
fn checkpointed_then_advanced(root: &Path) -> (ConversationManager, ConversationHistory) {
let manager = ConversationManager::new(root).expect("manager");
let mut state = driven_state(root);
let checkpointed = state.session.snapshot_conversation();
let events = state.session.drain_events(&checkpointed);
manager
.append_session_events(&checkpointed, &events)
.expect("append");
manager
.save_conversation(&checkpointed)
.expect("checkpoint");
state
.session
.append(ChatMessage::user("after the checkpoint"), fixed_ts());
state
.session
.append(ChatMessage::assistant("still here"), fixed_ts());
let latest = state.session.snapshot_conversation();
let tail = state.session.drain_events(&latest);
manager
.append_session_events(&latest, &tail)
.expect("append tail");
(manager, latest)
}
#[test]
fn resume_replays_the_tail_past_the_checkpoint() {
let root = temp_root("replay_tail");
let (manager, latest) = checkpointed_then_advanced(&root);
let resumed = manager.load_conversation(&latest.id).expect("resume");
assert_eq!(
json_of(&resumed),
json_of(&latest),
"resume must equal the live session, not the stale checkpoint"
);
assert!(
resumed
.messages()
.iter()
.any(|m| m.content == "after the checkpoint"),
"the events past the checkpoint must be replayed"
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn resume_ignores_a_checkpoint_it_cannot_place_in_the_log() {
let root = temp_root("untrusted_ckpt");
let (manager, latest) = checkpointed_then_advanced(&root);
let path = manager
.conversations_dir()
.join(format!("{}.json", latest.id));
let raw = std::fs::read_to_string(&path).expect("read checkpoint");
let mut value: serde_json::Value = serde_json::from_str(&raw).expect("parse");
assert!(
value
.as_object_mut()
.expect("object")
.remove("checkpoint_seq")
.is_some(),
"the checkpoint must have carried a watermark for this to prove anything"
);
std::fs::write(&path, serde_json::to_string(&value).unwrap()).expect("rewrite");
let resumed = manager.load_conversation(&latest.id).expect("resume");
assert_eq!(
json_of(&resumed),
json_of(&latest),
"a watermark-less checkpoint must fall back to a full fold"
);
let mut value: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
value.as_object_mut().unwrap().insert(
"checkpoint_seq".to_string(),
serde_json::Value::from(9_999_u64),
);
std::fs::write(&path, serde_json::to_string(&value).unwrap()).expect("rewrite");
let resumed = manager.load_conversation(&latest.id).expect("resume");
assert_eq!(
json_of(&resumed),
json_of(&latest),
"a checkpoint ahead of its log must fall back to a full fold"
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_deleted_checkpoint_costs_nothing_but_time() {
let root = temp_root("ckpt_deleted");
let (manager, latest) = checkpointed_then_advanced(&root);
std::fs::remove_file(
manager
.conversations_dir()
.join(format!("{}.json", latest.id)),
)
.expect("delete the checkpoint");
let resumed = manager.load_conversation(&latest.id).expect("resume");
assert_eq!(json_of(&resumed), json_of(&latest));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn continue_and_the_picker_see_a_session_with_no_checkpoint_yet() {
let root = temp_root("no_ckpt_yet");
let manager = ConversationManager::new(&root).expect("manager");
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager
.append_session_events(&snapshot, &events)
.expect("append");
assert!(
!manager
.conversations_dir()
.join(format!("{}.json", snapshot.id))
.exists(),
"this test is about the no-checkpoint case"
);
let last = manager
.load_last_conversation()
.expect("continue")
.expect("a session with a log is resumable");
assert_eq!(last.id, snapshot.id);
assert_eq!(json_of(&last), json_of(&snapshot));
let metas = manager.list_conversation_metas().expect("picker list");
assert_eq!(metas.len(), 1, "the picker must see it: {metas:?}");
assert_eq!(metas[0].id, snapshot.id);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_legacy_session_without_a_log_still_loads() {
let root = temp_root("legacy");
let manager = ConversationManager::new(&root).expect("manager");
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let _ = state.session.drain_events(&snapshot);
manager.save_conversation(&snapshot).expect("save");
std::fs::remove_file(
manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id)),
)
.ok();
assert!(
!manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id))
.exists()
);
let loaded = manager.load_conversation(&snapshot.id).expect("load");
assert_eq!(json_of(&loaded), json_of(&snapshot));
assert_eq!(
manager
.load_last_conversation()
.expect("continue")
.expect("legacy session is resumable")
.id,
snapshot.id
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_second_writer_diverts_this_process_to_a_conflict_sibling() {
let root = temp_root("log_conflict");
let manager = ConversationManager::new(&root).expect("manager");
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager
.append_session_events(&snapshot, &events)
.expect("our first append");
let shared = manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id));
let before = std::fs::read_to_string(&shared).expect("read shared");
let other = ConversationManager::new(&root).expect("other process");
let mut other_state = driven_state(&root);
other_state.session.conversation.id = snapshot.id.clone();
let other_snapshot = other_state.session.snapshot_conversation();
let other_events = other_state.session.drain_events(&other_snapshot);
other
.append_session_events(&other_snapshot, &other_events)
.expect("the other writer lands");
assert_ne!(
std::fs::read_to_string(&shared).unwrap().len(),
before.len(),
"the other writer must actually have grown the log"
);
state
.session
.append(ChatMessage::user("ours after the clash"), fixed_ts());
let latest = state.session.snapshot_conversation();
let ours = state.session.drain_events(&latest);
manager
.append_session_events(&latest, &ours)
.expect("diverting is not an error");
let siblings: Vec<PathBuf> = std::fs::read_dir(manager.conversations_dir())
.unwrap()
.flatten()
.map(|e| e.path())
.filter(|p| p.to_string_lossy().contains(".conflict.jsonl"))
.collect();
assert_eq!(siblings.len(), 1, "exactly one sibling: {siblings:?}");
let diverted = std::fs::read_to_string(&siblings[0]).unwrap();
assert!(
diverted.contains("ours after the clash"),
"our continuation must be preserved: {diverted}"
);
assert!(
!std::fs::read_to_string(&shared)
.unwrap()
.contains("ours after the clash"),
"and must NOT have interleaved into the shared log"
);
assert!(
diverted.contains("\"type\":\"started\""),
"the sibling must be foldable on its own: {diverted}"
);
let listed = manager.list_conversation_metas().expect("list");
assert_eq!(listed.len(), 1, "one session, not two: {listed:?}");
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_blocked_log_path_errors_instead_of_spinning() {
let root = temp_root("blocked_path");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
std::fs::create_dir_all(
manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id)),
)
.expect("plant a directory in the log's place");
assert!(
manager.append_session_events(&snapshot, &events).is_err(),
"append must fail on a blocked path"
);
assert!(
manager
.fold_conversation_from_log(&snapshot.id)
.expect("fold must not error out")
.is_none(),
"a non-file path folds to nothing"
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn newer_format_log_is_refused_not_misread() {
let root = temp_root("newer");
let manager = ConversationManager::new(&root).unwrap();
let id = "20260101_120000_001";
std::fs::write(
manager.conversations_dir().join(format!("{id}.jsonl")),
format!(
"{{\"v\":{},\"seq\":0,\"ts\":\"2026-07-02T12:00:00.123+00:00\",\"event\":{{\"type\":\"started\",\"session_id\":\"{id}\",\"project_path\":\"/p\",\"model_id\":\"m\",\"created_at\":\"2026-07-02T12:00:00.123+00:00\"}}}}\n",
SESSION_EVENT_FORMAT_VERSION + 1
),
)
.unwrap();
assert!(
manager.fold_conversation_from_log(id).unwrap().is_none(),
"a newer-format log must not be folded"
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn delete_conversation_removes_the_log_too() {
let root = temp_root("delete");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
manager.save_conversation(&snapshot).unwrap();
let log_path = manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id));
assert!(log_path.exists());
manager.delete_conversation(&snapshot.id).unwrap();
assert!(!log_path.exists(), "the log must cascade with the session");
let _ = std::fs::remove_dir_all(&root);
}
#[cfg(unix)]
#[test]
fn log_file_is_owner_only() {
use std::os::unix::fs::PermissionsExt;
let root = temp_root("perms");
let manager = ConversationManager::new(&root).unwrap();
let mut state = driven_state(&root);
let snapshot = state.session.snapshot_conversation();
let events = state.session.drain_events(&snapshot);
manager.append_session_events(&snapshot, &events).unwrap();
let mode = std::fs::metadata(
manager
.conversations_dir()
.join(format!("{}.jsonl", snapshot.id)),
)
.unwrap()
.permissions()
.mode()
& 0o777;
assert_eq!(mode, 0o600, "log must be created owner-only");
let _ = std::fs::remove_dir_all(&root);
}
}