use std::collections::HashMap;
use std::fs::{self, File, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
use anyhow::{Context, Result, anyhow, bail};
use chrono::{DateTime, FixedOffset, Local, Utc};
use serde::{Deserialize, Serialize};
use crate::llm::{Message, Role};
pub const FORMAT_VERSION: u32 = 1;
pub fn now() -> DateTime<FixedOffset> {
Local::now().fixed_offset()
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type", content = "data", rename_all = "snake_case")]
pub enum Record {
Session {
version: u32,
id: String,
created_at: DateTime<FixedOffset>,
},
Input {
id: String,
text: String,
recorded_at: DateTime<FixedOffset>,
},
Message(Message),
TurnEnd {
input_id: String,
response: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
outcome: Option<crate::goal::Outcome>,
#[serde(default, skip_serializing_if = "is_zero")]
history_calls: u32,
recorded_at: DateTime<FixedOffset>,
},
Replace {
messages: Vec<Message>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pending_position: Option<usize>,
#[serde(default, skip_serializing_if = "Option::is_none")]
summarized: Option<(u64, u64)>,
#[serde(default, skip_serializing_if = "Option::is_none")]
mode: Option<crate::config::CompactionMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
model: Option<String>,
recorded_at: DateTime<FixedOffset>,
},
Plan {
plan: crate::plan::Plan,
recorded_at: DateTime<FixedOffset>,
},
}
#[derive(Debug, Clone, PartialEq)]
pub struct PendingInput {
pub id: String,
pub text: String,
pub position: usize,
}
#[derive(Debug, Default)]
pub struct Restored {
pub conversation: Vec<Message>,
pub completed: HashMap<String, String>,
pub outcomes: HashMap<String, crate::goal::Outcome>,
pub pending_input: Option<PendingInput>,
pub plan: Option<crate::plan::Plan>,
pub history_available: bool,
pub history_hint_consumed: bool,
}
fn is_zero(n: &u32) -> bool {
*n == 0
}
pub struct SessionLog {
path: PathBuf,
file: File,
lines: u64,
}
pub fn spill_dir_for(dir: &Path, id: &str) -> PathBuf {
dir.join(format!("{id}.spill"))
}
pub fn default_dir() -> PathBuf {
crate::config::app_dir(&dirs::data_local_dir().unwrap_or_else(std::env::temp_dir)).join("sessions")
}
pub fn new_session_id() -> String {
format!("sess-{}-{:08x}", Utc::now().format("%Y%m%dT%H%M%S"), fastrand::u32(..))
}
pub fn validate_id(id: &str) -> Result<()> {
let valid = !id.is_empty()
&& id.len() <= 128
&& !id.starts_with('.')
&& id.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.'));
if !valid {
bail!("invalid session id {id:?}: use 1-128 of [A-Za-z0-9._-], not starting with '.'");
}
Ok(())
}
fn path_for(dir: &Path, id: &str) -> PathBuf {
dir.join(format!("{id}.jsonl"))
}
fn encode(record: &Record) -> Result<Vec<u8>> {
let mut line = serde_json::to_vec(record).context("encode session record")?;
line.push(b'\n');
Ok(line)
}
impl SessionLog {
pub fn create(dir: &Path, id: &str) -> Result<Self> {
validate_id(id)?;
fs::create_dir_all(dir).with_context(|| format!("create session dir {}", dir.display()))?;
let path = path_for(dir, id);
let mut file = OpenOptions::new()
.create_new(true)
.append(true)
.open(&path)
.with_context(|| format!("create session log {}", path.display()))?;
file.write_all(&encode(&Record::Session { version: FORMAT_VERSION, id: id.to_string(), created_at: now() })?)?;
file.sync_data()?;
Ok(Self { path, file, lines: 1 })
}
pub fn open(dir: &Path, id: &str) -> Result<(Self, Restored)> {
validate_id(id)?;
let path = path_for(dir, id);
let bytes = fs::read(&path).with_context(|| format!("read session log {}", path.display()))?;
let committed = bytes
.iter()
.rposition(|&b| b == b'\n')
.map(|i| i + 1)
.ok_or_else(|| anyhow!("session log {} has no committed records", path.display()))?;
let restored =
decode(&bytes[..committed], id).with_context(|| format!("load session log {}", path.display()))?;
let file = OpenOptions::new()
.append(true)
.open(&path)
.with_context(|| format!("open session log {}", path.display()))?;
if committed < bytes.len() {
file.set_len(committed as u64)?;
}
let lines = bytes[..committed].iter().filter(|&&b| b == b'\n').count() as u64;
Ok((Self { path, file, lines }, restored))
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn append(&mut self, record: &Record) -> Result<u64> {
self.file
.write_all(&encode(record)?)
.with_context(|| format!("append to session log {}", self.path.display()))?;
if matches!(record, Record::TurnEnd { .. } | Record::Replace { .. }) {
self.file.sync_data()?;
}
self.lines += 1;
Ok(self.lines)
}
}
pub fn read_records(dir: &Path, id: &str) -> Result<Vec<Record>> {
validate_id(id)?;
read_records_at(&path_for(dir, id), Some(id))
}
pub fn read_records_at(path: &Path, expected_id: Option<&str>) -> Result<Vec<Record>> {
let bytes = fs::read(path).with_context(|| format!("read session log {}", path.display()))?;
let committed = bytes
.iter()
.rposition(|&b| b == b'\n')
.map(|i| i + 1)
.ok_or_else(|| anyhow!("session log {} has no committed records", path.display()))?;
let mut records = Vec::new();
let mut index = 0usize;
for (line_no, line) in bytes[..committed].split(|&b| b == b'\n').enumerate() {
if line.is_empty() {
continue;
}
let record: Record = serde_json::from_slice(line)
.with_context(|| format!("decode record {} of {}", index + 1, path.display()))?;
match &record {
Record::Session { version, id, .. } => {
if index != 0 {
bail!("session record must be first (found at record {})", index + 1);
}
if *version != FORMAT_VERSION {
bail!("unsupported session format version {version} (this build supports {FORMAT_VERSION})");
}
if let Some(expected) = expected_id
&& id != expected
{
bail!("session log header id {id:?} does not match {expected:?}");
}
}
_ if index == 0 => bail!("session log {} does not start with a session header", path.display()),
_ => {}
}
let mut record = record;
if let Record::Message(message) = &mut record {
message.log_line = Some(line_no as u64 + 1);
}
records.push(record);
index += 1;
}
if records.is_empty() {
bail!("session log {} has no decodable records", path.display());
}
Ok(records)
}
fn same_content(a: &Message, b: &Message) -> bool {
a.role == b.role
&& a.content == b.content
&& a.tool_calls == b.tool_calls
&& a.tool_call_id == b.tool_call_id
&& a.name == b.name
&& a.is_error == b.is_error
&& a.thinking_blocks == b.thinking_blocks
}
fn backfill_log_lines(messages: &mut [Message], originals: &[Message], summarized: Option<(u64, u64)>) {
let floor = summarized.map_or(0, |(_, end)| end);
let mut claimed = vec![false; originals.len()];
for message in messages.iter_mut().filter(|m| m.log_line.is_none() && m.role != Role::System) {
if let Some((index, original)) = originals
.iter()
.enumerate()
.find(|(i, o)| !claimed[*i] && o.log_line.is_some_and(|line| line > floor) && same_content(o, message))
{
message.log_line = original.log_line;
claimed[index] = true;
}
}
}
fn decode(bytes: &[u8], expected_id: &str) -> Result<Restored> {
let mut restored = Restored::default();
let mut originals: Vec<Message> = Vec::new();
let mut index = 0usize;
for (line_no, line) in bytes.split(|&b| b == b'\n').enumerate() {
if line.is_empty() {
continue;
}
let record: Record = serde_json::from_slice(line).with_context(|| format!("decode record {}", index + 1))?;
match record {
Record::Session { version, id, .. } => {
if index != 0 {
bail!("session record must be first (found at record {})", index + 1);
}
if version != FORMAT_VERSION {
bail!("unsupported session format version {version} (this build supports {FORMAT_VERSION})");
}
if id != expected_id {
bail!("session log header id {id:?} does not match {expected_id:?}");
}
}
_ if index == 0 => bail!("session log does not start with a session header"),
Record::Input { id, text, .. } => {
restored.pending_input = Some(PendingInput { id, text, position: restored.conversation.len() });
}
Record::Message(mut message) => {
message.log_line = Some(line_no as u64 + 1);
if crate::history::consumes_hint(&message) {
restored.history_hint_consumed = true;
}
originals.push(message.clone());
restored.conversation.push(message);
}
Record::TurnEnd { input_id, response, outcome, .. } => {
if restored.pending_input.as_ref().is_some_and(|p| p.id == input_id) {
restored.pending_input = None;
}
match outcome {
Some(outcome) => restored.outcomes.insert(input_id.clone(), outcome),
None => restored.outcomes.remove(&input_id),
};
restored.completed.insert(input_id, response);
}
Record::Replace { mut messages, pending_position, summarized, mode, .. } => {
backfill_log_lines(&mut messages, &originals, summarized);
restored.conversation = messages;
restored.pending_input = match (restored.pending_input.take(), pending_position) {
(Some(pending), Some(position)) => Some(PendingInput { position, ..pending }),
_ => None,
};
restored.history_available = matches!(mode, Some(crate::config::CompactionMode::Smart));
restored.history_hint_consumed = false;
}
Record::Plan { plan, .. } => restored.plan = Some(plan),
}
index += 1;
}
Ok(restored)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_stale_log_line_on_a_direct_record_is_renumbered_to_its_physical_line() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("s.jsonl");
let stale: Record =
serde_json::from_str(r#"{"type":"message","data":{"role":"user","content":"hi","log_line":99}}"#).unwrap();
assert_eq!(stale, Record::Message(Message { log_line: Some(99), ..Message::user("hi") }));
let mut log = String::new();
for record in
[Record::Session { version: FORMAT_VERSION, id: "s".into(), created_at: now() }, input("i"), stale]
{
log.push_str(&serde_json::to_string(&record).unwrap());
log.push('\n');
}
std::fs::write(&path, log).unwrap();
let records = read_records_at(&path, Some("s")).unwrap();
let Record::Message(message) = &records[2] else { panic!("a message record") };
assert_eq!(message.log_line, Some(3), "the physical line wins over the serialized field");
let (_, restored) = SessionLog::open(dir.path(), "s").unwrap();
assert_eq!(restored.conversation[0].log_line, Some(3));
}
#[test]
fn reads_utc_records_and_unstamped_messages_from_older_logs() {
let input: Record = serde_json::from_str(
r#"{"type":"input","data":{"id":"i","text":"hi","recorded_at":"2026-01-01T00:00:00Z"}}"#,
)
.unwrap();
let Record::Input { recorded_at, .. } = input else { panic!() };
assert_eq!(recorded_at.to_rfc3339(), "2026-01-01T00:00:00+00:00");
let message: Record =
serde_json::from_str(r#"{"type":"message","data":{"role":"user","content":"hi"}}"#).unwrap();
assert_eq!(message, Record::Message(Message::user("hi")));
let now = serde_json::to_string(&now()).unwrap();
assert!(now.contains('+') || now.contains("-0") || now.contains("-1"), "{now}");
}
fn input(id: &str) -> Record {
Record::Input { id: id.into(), text: "hi".into(), recorded_at: now() }
}
fn turn_end(id: &str) -> Record {
Record::TurnEnd {
input_id: id.into(),
response: "hello".into(),
outcome: None,
history_calls: 0,
recorded_at: now(),
}
}
#[test]
fn round_trips_and_tracks_inputs() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s1").unwrap();
log.append(&Record::Message(Message::system("sys"))).unwrap();
log.append(&input("in-1")).unwrap();
log.append(&Record::Message(Message::user("hi"))).unwrap();
log.append(&Record::Message(Message::assistant("hello"))).unwrap();
log.append(&turn_end("in-1")).unwrap();
log.append(&input("in-2")).unwrap();
log.append(&Record::Message(Message::user("again"))).unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s1").unwrap();
assert_eq!(restored.conversation.len(), 4);
assert_eq!(restored.completed.get("in-1").map(String::as_str), Some("hello"));
assert_eq!(restored.pending_input, Some(PendingInput { id: "in-2".into(), text: "hi".into(), position: 3 }));
assert!(SessionLog::create(dir.path(), "s1").is_err(), "create must not clobber");
}
#[test]
fn replace_resets_conversation() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s2").unwrap();
log.append(&Record::Message(Message::user("a"))).unwrap();
log.append(&Record::Replace {
messages: vec![Message::system("new")],
pending_position: None,
summarized: None,
mode: None,
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s2").unwrap();
assert_eq!(restored.conversation, vec![Message::system("new")]);
}
#[test]
fn replace_backfills_log_lines_for_legacy_retained_messages() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s5").unwrap();
log.append(&Record::Message(Message::system("sys"))).unwrap(); log.append(&Record::Message(Message::user("first"))).unwrap(); log.append(&Record::Message(Message::assistant("reply"))).unwrap(); log.append(&Record::Replace {
messages: vec![
Message::system("sys"),
Message::user("first"),
Message::assistant("reply"),
Message::assistant("summary"),
],
pending_position: None,
summarized: None,
mode: None,
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s5").unwrap();
let lines: Vec<Option<u64>> = restored.conversation.iter().map(|m| m.log_line).collect();
assert_eq!(lines, vec![None, Some(3), Some(4), None]);
}
#[test]
fn replace_can_keep_the_pending_input() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s3").unwrap();
log.append(&Record::Input { id: "in-1".into(), text: "go".into(), recorded_at: now() }).unwrap();
log.append(&Record::Message(Message::user("go"))).unwrap();
let messages = vec![Message::system("sys"), Message::user("summary"), Message::user("go")];
log.append(&Record::Replace {
messages: messages.clone(),
pending_position: Some(2),
summarized: None,
mode: None,
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s3").unwrap();
let mut expected = messages.clone();
expected[2].log_line = Some(3);
assert_eq!(restored.conversation, expected);
assert_eq!(restored.pending_input, Some(PendingInput { id: "in-1".into(), text: "go".into(), position: 2 }));
}
#[test]
fn backfill_uses_the_summarized_range_to_disambiguate_duplicates() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s6").unwrap();
log.append(&Record::Message(Message::system("sys"))).unwrap(); log.append(&Record::Message(Message::user("same"))).unwrap(); log.append(&Record::Message(Message::assistant("x"))).unwrap(); log.append(&Record::Message(Message::user("same"))).unwrap(); log.append(&Record::Replace {
messages: vec![Message::user("summary"), Message::user("same")],
pending_position: None,
summarized: Some((2, 4)),
mode: Some(crate::config::CompactionMode::Smart),
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s6").unwrap();
assert_eq!(restored.conversation[1].log_line, Some(5));
assert!(restored.history_available, "a smart replace offers the history tools");
}
#[test]
fn blank_lines_do_not_shift_message_ids() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("s7.jsonl");
let header =
serde_json::to_string(&Record::Session { version: FORMAT_VERSION, id: "s7".into(), created_at: now() })
.unwrap();
let user = serde_json::to_string(&Record::Message(Message::user("hi"))).unwrap();
let answer = serde_json::to_string(&Record::Message(Message::assistant("hello"))).unwrap();
std::fs::write(&path, format!("{header}\n{user}\n\n{answer}\n")).unwrap();
let (_, restored) = SessionLog::open(dir.path(), "s7").unwrap();
let lines: Vec<Option<u64>> = restored.conversation.iter().map(|m| m.log_line).collect();
assert_eq!(lines, vec![Some(2), Some(4)], "the assistant message stays on physical line 4");
let (mut log, _) = SessionLog::open(dir.path(), "s7").unwrap();
let line = log.append(&Record::Message(Message::user("again"))).unwrap();
assert_eq!(line, 5, "the append lands on physical line 5, past the blank line 3");
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s7").unwrap();
let lines: Vec<Option<u64>> = restored.conversation.iter().map(|m| m.log_line).collect();
assert_eq!(lines, vec![Some(2), Some(4), Some(5)]);
let records = read_records(dir.path(), "s7").unwrap();
let ids: Vec<Option<u64>> = records
.iter()
.filter_map(|r| match r {
Record::Message(m) => Some(m.log_line),
_ => None,
})
.collect();
assert_eq!(ids, vec![Some(2), Some(4), Some(5)]);
}
#[test]
fn standard_replace_does_not_offer_history_tools() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s7").unwrap();
log.append(&Record::Message(Message::system("sys"))).unwrap();
log.append(&Record::Replace {
messages: vec![Message::user("summary")],
pending_position: None,
summarized: Some((2, 2)),
mode: Some(crate::config::CompactionMode::Standard),
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s7").unwrap();
assert!(!restored.history_available, "a standard replace drops the history tools");
}
#[test]
fn hint_consumed_tracks_messages_after_the_latest_replace() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s7b").unwrap();
log.append(&Record::Message(Message::system("sys"))).unwrap();
log.append(&Record::Replace {
messages: vec![Message::user("summary")],
pending_position: None,
summarized: Some((2, 2)),
mode: Some(crate::config::CompactionMode::Smart),
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s7b").unwrap();
assert!(restored.history_available && !restored.history_hint_consumed);
let mut log = SessionLog::open(dir.path(), "s7b").unwrap().0;
let hinted = format!("Exit code: 1\n\n{}", crate::history::FAILED_TOOL_HINT);
log.append(&Record::Message(Message::tool_error("c1", "bash", &hinted))).unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s7b").unwrap();
assert!(restored.history_hint_consumed, "a hint emitted after the replace is consumed");
}
#[test]
fn backfill_distinguishes_messages_by_thinking_blocks() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s8").unwrap();
let think_a =
Message { thinking_blocks: vec![serde_json::json!({"thinking": "a"})], ..Message::assistant("reply") };
let think_b =
Message { thinking_blocks: vec![serde_json::json!({"thinking": "b"})], ..Message::assistant("reply") };
log.append(&Record::Message(Message::system("sys"))).unwrap(); log.append(&Record::Message(think_a.clone())).unwrap(); log.append(&Record::Message(think_b.clone())).unwrap(); log.append(&Record::Replace {
messages: vec![Message::user("summary"), Message { log_line: None, ..think_b.clone() }],
pending_position: None,
summarized: None,
mode: None,
model: None,
recorded_at: now(),
})
.unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s8").unwrap();
assert_eq!(restored.conversation[1].log_line, Some(4));
}
#[test]
fn discards_torn_tail_and_appends_cleanly() {
let dir = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(dir.path(), "s3").unwrap();
log.append(&Record::Message(Message::user("a"))).unwrap();
drop(log);
let path = path_for(dir.path(), "s3");
let mut file = OpenOptions::new().append(true).open(&path).unwrap();
file.write_all(br#"{"type":"message","data":{"role":"us"#).unwrap();
drop(file);
let (mut log, restored) = SessionLog::open(dir.path(), "s3").unwrap();
assert_eq!(restored.conversation.len(), 1);
log.append(&Record::Message(Message::user("b"))).unwrap();
drop(log);
let (_, restored) = SessionLog::open(dir.path(), "s3").unwrap();
assert_eq!(restored.conversation.len(), 2);
}
#[test]
fn rejects_unsupported_versions_and_bad_ids() {
let dir = tempfile::tempdir().unwrap();
fs::write(
path_for(dir.path(), "future"),
"{\"type\":\"session\",\"data\":{\"version\":99,\"id\":\"future\",\"created_at\":\"2026-01-01T00:00:00Z\"}}\n",
)
.unwrap();
let err = SessionLog::open(dir.path(), "future").err().unwrap();
assert!(format!("{err:#}").contains("unsupported session format version 99"), "{err:#}");
assert!(validate_id("../etc/passwd").is_err());
assert!(validate_id(".hidden").is_err());
assert!(validate_id(&new_session_id()).is_ok());
}
#[test]
fn golden_format_is_stable() {
let record = Record::Message(Message::tool_result("c1", "bash", "ok"));
assert_eq!(
serde_json::to_string(&record).unwrap(),
r#"{"type":"message","data":{"role":"tool","content":"ok","tool_call_id":"c1","name":"bash"}}"#
);
}
#[test]
fn read_records_rejects_a_header_id_that_differs_from_the_requested_id() {
let dir = tempfile::tempdir().unwrap();
let path = path_for(dir.path(), "renamed");
fs::write(
&path,
"{\"type\":\"session\",\"data\":{\"version\":1,\"id\":\"real\",\"created_at\":\"2026-01-01T00:00:00Z\"}}\n",
)
.unwrap();
let err = read_records(dir.path(), "renamed").err().unwrap();
assert!(format!("{err:#}").contains("does not match"), "{err:#}");
assert!(read_records_at(&path, None).is_ok());
assert!(read_records_at(&path, Some("real")).is_ok());
}
}