use std::io::Write;
use std::path::{Path, PathBuf};
use locode_protocol::{Event, Message};
use serde::{Deserialize, Serialize};
use crate::session_dirs::encode_cwd_dirname;
pub const TRACE_SCHEMA_VERSION: u32 = 1;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionMeta {
pub schema_version: u32,
pub session_id: String,
#[serde(default = "default_kind")]
pub kind: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub group: Option<String>,
pub cwd: PathBuf,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub git: Option<GitMeta>,
pub cli_version: String,
pub harness: String,
pub api_schema: String,
pub model: String,
}
fn default_kind() -> String {
"main".to_string()
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct GitMeta {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub root: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub head: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub remote: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct TraceExtras {
pub cli_version: String,
pub git: Option<GitMeta>,
pub kind: Option<String>,
pub parent_id: Option<String>,
pub group: Option<String>,
}
#[derive(Debug)]
pub struct TraceWriter {
sessions_root: PathBuf,
extras: TraceExtras,
file: Option<std::fs::File>,
path: Option<PathBuf>,
error: Option<String>,
resumed: bool,
}
impl TraceWriter {
#[must_use]
pub fn new(sessions_root: PathBuf, extras: TraceExtras) -> Self {
Self {
sessions_root,
extras,
file: None,
path: None,
error: None,
resumed: false,
}
}
pub fn resume(path: PathBuf, sessions_root: PathBuf) -> Result<Self, String> {
let file = open_append_private(&path)?;
Ok(Self {
sessions_root,
extras: TraceExtras::default(),
file: Some(file),
path: Some(path),
error: None,
resumed: true,
})
}
pub fn on_event(&mut self, event: &Event) {
if self.error.is_some() {
return;
}
let result = match event {
Event::Init {
session_id,
harness,
api_schema,
model,
cwd,
preamble,
..
} if !self.resumed => {
self.on_init(session_id, harness, api_schema, model, cwd, preamble)
}
Event::Message { message } => self.append_message(message),
Event::Result { report } => {
if self.file.is_some() {
self.append_record("usage", &report.usage)
} else {
Ok(())
}
}
_ => Ok(()),
};
if let Err(e) = result {
self.error = Some(e);
self.file = None;
}
}
#[must_use]
pub fn path(&self) -> Option<&Path> {
self.path.as_deref()
}
pub fn take_error(&mut self) -> Option<String> {
self.error.take()
}
fn on_init(
&mut self,
session_id: &str,
harness: &str,
api_schema: &str,
model: &str,
cwd: &str,
preamble: &[Message],
) -> Result<(), String> {
let cwd_path = PathBuf::from(cwd);
let dirname = encode_cwd_dirname(&cwd_path);
let dir = self.sessions_root.join(&dirname);
create_dir_private(&dir)?;
if !dirname.starts_with('+') {
let sidecar = dir.join(".cwd");
if !sidecar.exists() {
std::fs::write(&sidecar, cwd).map_err(|e| format!("write .cwd sidecar: {e}"))?;
}
}
let filename = format!("rollout-{}-{session_id}.jsonl", filename_timestamp());
let path = dir.join(&filename);
let file = open_append_private(&path)?;
self.file = Some(file);
self.path = Some(path);
let meta = SessionMeta {
schema_version: TRACE_SCHEMA_VERSION,
session_id: session_id.to_string(),
kind: self
.extras
.kind
.clone()
.unwrap_or_else(|| "main".to_string()),
parent_id: self.extras.parent_id.clone(),
group: self.extras.group.clone(),
cwd: cwd_path,
git: self.extras.git.clone(),
cli_version: self.extras.cli_version.clone(),
harness: harness.to_string(),
api_schema: api_schema.to_string(),
model: model.to_string(),
};
self.append_record("session_meta", &meta)?;
for message in preamble {
self.append_record("message", message)?;
}
Ok(())
}
fn append_message(&mut self, message: &Message) -> Result<(), String> {
if self.file.is_none() {
return Ok(()); }
self.append_record("message", message)
}
fn append_record(&mut self, record_type: &str, payload: &impl Serialize) -> Result<(), String> {
let Some(file) = self.file.as_mut() else {
return Ok(());
};
let line = serde_json::json!({
"timestamp": now_rfc3339_millis(),
"type": record_type,
"payload": payload,
});
let mut buf = serde_json::to_string(&line).map_err(|e| format!("serialize: {e}"))?;
buf.push('\n');
file.write_all(buf.as_bytes())
.and_then(|()| file.flush())
.map_err(|e| format!("append trace record: {e}"))
}
}
#[derive(Debug)]
pub struct RolloutContents {
pub meta: SessionMeta,
pub history: Vec<Message>,
pub last_usage: Option<locode_protocol::Usage>,
}
pub fn read_rollout(path: &Path) -> Result<RolloutContents, String> {
let text =
std::fs::read_to_string(path).map_err(|e| format!("read {}: {e}", path.display()))?;
let mut lines = text.lines();
let first = lines
.next()
.ok_or_else(|| format!("{}: empty rollout", path.display()))?;
let first: serde_json::Value = serde_json::from_str(first)
.map_err(|e| format!("{}: line 1 unparsable: {e}", path.display()))?;
if first.get("type").and_then(serde_json::Value::as_str) != Some("session_meta") {
return Err(format!("{}: line 1 is not session_meta", path.display()));
}
let meta: SessionMeta = serde_json::from_value(first["payload"].clone())
.map_err(|e| format!("{}: session_meta invalid: {e}", path.display()))?;
let mut history: Vec<Message> = Vec::new();
let mut last_usage: Option<locode_protocol::Usage> = None;
for line in lines {
let Ok(value) = serde_json::from_str::<serde_json::Value>(line) else {
continue;
};
match value.get("type").and_then(serde_json::Value::as_str) {
Some("message") => {
if let Ok(message) = serde_json::from_value::<Message>(value["payload"].clone()) {
history.push(message);
}
}
Some("compacted") => {
if let Some(replacement) = value["payload"].get("replacement_history")
&& let Ok(messages) =
serde_json::from_value::<Vec<Message>>(replacement.clone())
{
history = messages;
}
}
Some("usage") => {
if let Ok(usage) =
serde_json::from_value::<locode_protocol::Usage>(value["payload"].clone())
{
last_usage = Some(usage);
}
}
_ => {}
}
}
Ok(RolloutContents {
meta,
history,
last_usage,
})
}
#[must_use]
pub fn find_latest_rollout(sessions_root: &Path, cwd: &Path) -> Option<PathBuf> {
let dir = sessions_root.join(encode_cwd_dirname(cwd));
let mut names: Vec<String> = rollout_names(&dir);
names.sort_unstable_by(|a, b| b.cmp(a)); names
.into_iter()
.map(|n| dir.join(n))
.find(|path| read_rollout(path).is_ok_and(|c| c.meta.kind == "main"))
}
#[must_use]
pub fn find_rollout_by_id(sessions_root: &Path, cwd: &Path, id: &str) -> Option<PathBuf> {
let suffix = format!("-{id}.jsonl");
let scoped = sessions_root.join(encode_cwd_dirname(cwd));
if let Some(hit) = dir_hit(&scoped, &suffix) {
return Some(hit);
}
let entries = std::fs::read_dir(sessions_root).ok()?;
for entry in entries.flatten() {
let dir = entry.path();
if dir == scoped || !dir.is_dir() {
continue;
}
if let Some(hit) = dir_hit(&dir, &suffix) {
return Some(hit);
}
}
None
}
fn rollout_names(dir: &Path) -> Vec<String> {
let Ok(entries) = std::fs::read_dir(dir) else {
return Vec::new();
};
entries
.flatten()
.filter_map(|e| e.file_name().into_string().ok())
.filter(|n| {
n.starts_with("rollout-")
&& std::path::Path::new(n)
.extension()
.is_some_and(|ext| ext.eq_ignore_ascii_case("jsonl"))
})
.collect()
}
fn dir_hit(dir: &Path, suffix: &str) -> Option<PathBuf> {
rollout_names(dir)
.into_iter()
.find(|n| n.ends_with(suffix))
.map(|n| dir.join(n))
}
pub(crate) fn open_append_private(path: &Path) -> Result<std::fs::File, String> {
let mut options = std::fs::OpenOptions::new();
options.create(true).append(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
options.mode(0o600);
}
let mut file = options
.open(path)
.map_err(|e| format!("open {}: {e}", path.display()))?;
heal_torn_tail(path, &mut file)?;
Ok(file)
}
fn heal_torn_tail(path: &Path, file: &mut std::fs::File) -> Result<(), String> {
use std::io::{Read, Seek, SeekFrom};
let len = file
.metadata()
.map_err(|e| format!("stat {}: {e}", path.display()))?
.len();
if len == 0 {
return Ok(());
}
let mut reader =
std::fs::File::open(path).map_err(|e| format!("open {}: {e}", path.display()))?;
reader
.seek(SeekFrom::End(-1))
.map_err(|e| format!("seek {}: {e}", path.display()))?;
let mut last = [0_u8; 1];
reader
.read_exact(&mut last)
.map_err(|e| format!("read {}: {e}", path.display()))?;
if last[0] != b'\n' {
file.write_all(b"\n")
.map_err(|e| format!("heal {}: {e}", path.display()))?;
}
Ok(())
}
pub(crate) fn create_dir_private(dir: &Path) -> Result<(), String> {
#[cfg(unix)]
{
use std::os::unix::fs::DirBuilderExt;
std::fs::DirBuilder::new()
.recursive(true)
.mode(0o700)
.create(dir)
.map_err(|e| format!("create {}: {e}", dir.display()))
}
#[cfg(not(unix))]
{
std::fs::create_dir_all(dir).map_err(|e| format!("create {}: {e}", dir.display()))
}
}
fn now_rfc3339_millis() -> String {
chrono::Utc::now()
.format("%Y-%m-%dT%H:%M:%S%.3fZ")
.to_string()
}
fn filename_timestamp() -> String {
chrono::Utc::now().format("%Y-%m-%dT%H-%M-%S").to_string()
}
#[cfg(test)]
mod tests {
use super::*;
use locode_protocol::{ContentBlock, Role};
use serde_json::Value;
fn init_event(session_id: &str, cwd: &Path) -> Event {
Event::Init {
session_id: session_id.to_string(),
harness: "codex".to_string(),
api_schema: "mock".to_string(),
model: "mock-1".to_string(),
cwd: cwd.to_string_lossy().into_owned(),
max_turns: None,
preamble: vec![Message {
role: Role::System,
content: vec![ContentBlock::Text {
text: "base prompt".to_string(),
}],
}],
tools: vec![],
}
}
fn message_event(role: Role, text: &str) -> Event {
Event::Message {
message: Message {
role,
content: vec![ContentBlock::Text {
text: text.to_string(),
}],
},
}
}
fn read_lines(path: &Path) -> Vec<Value> {
std::fs::read_to_string(path)
.unwrap()
.lines()
.map(|l| serde_json::from_str(l).unwrap())
.collect()
}
#[test]
fn writes_header_preamble_and_messages() {
let root = tempfile::tempdir().unwrap();
let cwd = tempfile::tempdir().unwrap();
let cwd = std::fs::canonicalize(cwd.path()).unwrap();
let mut writer = TraceWriter::new(
root.path().join("sessions"),
TraceExtras {
cli_version: "0.1.9".to_string(),
git: Some(GitMeta {
branch: Some("main".to_string()),
..Default::default()
}),
..Default::default()
},
);
writer.on_event(&init_event("sess-1", &cwd));
writer.on_event(&message_event(Role::User, "hi"));
writer.on_event(&Event::MessageDelta {
text: "ignored".to_string(),
});
writer.on_event(&message_event(Role::Assistant, "hello"));
assert!(writer.take_error().is_none());
let path = writer.path().unwrap().to_path_buf();
assert_eq!(
path.parent()
.unwrap()
.file_name()
.unwrap()
.to_str()
.unwrap(),
encode_cwd_dirname(&cwd)
);
assert!(
path.file_name()
.unwrap()
.to_str()
.unwrap()
.starts_with("rollout-")
);
assert!(
path.file_name()
.unwrap()
.to_str()
.unwrap()
.ends_with("-sess-1.jsonl")
);
let lines = read_lines(&path);
assert_eq!(lines.len(), 4, "meta + preamble + 2 messages");
assert_eq!(lines[0]["type"], "session_meta");
let meta = &lines[0]["payload"];
assert_eq!(meta["schema_version"], 1);
assert_eq!(meta["session_id"], "sess-1");
assert_eq!(meta["kind"], "main");
assert_eq!(meta["harness"], "codex");
assert_eq!(meta["git"]["branch"], "main");
assert!(meta.get("parent_id").is_none(), "absent when None");
assert_eq!(lines[1]["type"], "message");
assert_eq!(lines[1]["payload"]["role"], "system");
assert_eq!(lines[2]["payload"]["role"], "user");
assert_eq!(lines[3]["payload"]["role"], "assistant");
for line in &lines {
assert!(line["timestamp"].as_str().unwrap().ends_with('Z'));
}
}
#[cfg(unix)]
#[test]
fn modes_are_private() {
use std::os::unix::fs::PermissionsExt;
let root = tempfile::tempdir().unwrap();
let cwd = tempfile::tempdir().unwrap();
let cwd = std::fs::canonicalize(cwd.path()).unwrap();
let mut writer = TraceWriter::new(root.path().join("sessions"), TraceExtras::default());
writer.on_event(&init_event("sess-2", &cwd));
let path = writer.path().unwrap();
let file_mode = std::fs::metadata(path).unwrap().permissions().mode() & 0o777;
let dir_mode = std::fs::metadata(path.parent().unwrap())
.unwrap()
.permissions()
.mode()
& 0o777;
assert_eq!(file_mode, 0o600, "rollout is 0600");
assert_eq!(dir_mode, 0o700, "session dir is 0700");
}
#[test]
fn torn_tail_is_healed_on_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("rollout-x.jsonl");
std::fs::write(&path, "{\"ok\":1}\n{\"torn\":").unwrap();
let mut file = open_append_private(&path).unwrap();
file.write_all(b"{\"next\":2}\n").unwrap();
drop(file);
let text = std::fs::read_to_string(&path).unwrap();
assert_eq!(text, "{\"ok\":1}\n{\"torn\":\n{\"next\":2}\n");
let parsed: Vec<_> = text
.lines()
.filter(|l| serde_json::from_str::<Value>(l).is_ok())
.collect();
assert_eq!(parsed.len(), 2);
}
#[test]
fn no_init_means_no_file_and_errors_disable() {
let root = tempfile::tempdir().unwrap();
let mut writer = TraceWriter::new(root.path().join("sessions"), TraceExtras::default());
writer.on_event(&message_event(Role::User, "before init"));
assert!(writer.path().is_none(), "nothing written before Init");
let mut writer = TraceWriter::new(PathBuf::from("/dev/null/nope"), TraceExtras::default());
let cwd = std::env::temp_dir();
writer.on_event(&init_event("sess-3", &cwd));
assert!(writer.take_error().is_some());
writer.on_event(&message_event(Role::User, "x"));
}
#[test]
fn read_rollout_replays_and_tolerates() {
let root = tempfile::tempdir().unwrap();
let cwd = tempfile::tempdir().unwrap();
let cwd = std::fs::canonicalize(cwd.path()).unwrap();
let mut writer = TraceWriter::new(root.path().join("sessions"), TraceExtras::default());
writer.on_event(&init_event("sess-r", &cwd));
writer.on_event(&message_event(Role::User, "hi"));
writer.on_event(&message_event(Role::Assistant, "yo"));
let path = writer.path().unwrap().to_path_buf();
{
let mut file = open_append_private(&path).unwrap();
file.write_all(
b"{\"timestamp\":\"x\",\"type\":\"future_thing\",\"payload\":{\"n\":1}}\n",
)
.unwrap();
file.write_all(b"{\"torn").unwrap();
}
let contents = read_rollout(&path).unwrap();
assert_eq!(contents.meta.session_id, "sess-r");
assert_eq!(contents.history.len(), 3);
assert_eq!(contents.history[0].role, Role::System);
assert_eq!(contents.history[2].role, Role::Assistant);
}
#[test]
fn read_rollout_folds_compacted() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("rollout-t-sess-c.jsonl");
let meta = serde_json::json!({"timestamp":"t","type":"session_meta","payload":{
"schema_version":1,"session_id":"sess-c","kind":"main","cwd":"/x",
"cli_version":"0","harness":"grok","api_schema":"mock","model":"m"}});
let msg = |role: &str, text: &str| {
serde_json::json!({"timestamp":"t","type":"message","payload":
{"role":role,"content":[{"type":"text","text":text}]}})
};
let compacted = serde_json::json!({"timestamp":"t","type":"compacted","payload":{
"summary":"s","replacement_history":[
{"role":"user","content":[{"type":"text","text":"summary of the past"}]}]}});
let lines = [
meta.to_string(),
msg("user", "old-1").to_string(),
msg("assistant", "old-2").to_string(),
compacted.to_string(),
msg("user", "after").to_string(),
];
std::fs::write(&path, lines.join("\n") + "\n").unwrap();
let contents = read_rollout(&path).unwrap();
assert_eq!(contents.history.len(), 2);
assert_eq!(contents.history[1].role, Role::User);
}
#[test]
fn resumed_writer_appends_without_a_new_header() {
let root = tempfile::tempdir().unwrap();
let cwd = tempfile::tempdir().unwrap();
let cwd = std::fs::canonicalize(cwd.path()).unwrap();
let mut writer = TraceWriter::new(root.path().join("sessions"), TraceExtras::default());
writer.on_event(&init_event("sess-a", &cwd));
writer.on_event(&message_event(Role::User, "first run"));
let path = writer.path().unwrap().to_path_buf();
let before = read_lines(&path).len();
drop(writer);
let mut resumed = TraceWriter::resume(path.clone(), root.path().join("sessions")).unwrap();
resumed.on_event(&init_event("sess-a", &cwd)); resumed.on_event(&message_event(Role::User, "second run"));
assert!(resumed.take_error().is_none());
let lines = read_lines(&path);
assert_eq!(lines.len(), before + 1, "exactly one new message line");
assert_eq!(
lines.iter().filter(|l| l["type"] == "session_meta").count(),
1,
"one header, ever"
);
}
#[test]
fn find_latest_scopes_by_cwd_and_kind() {
let root = tempfile::tempdir().unwrap();
let sessions = root.path().join("sessions");
let cwd_a = tempfile::tempdir().unwrap();
let cwd_a = std::fs::canonicalize(cwd_a.path()).unwrap();
let cwd_b = tempfile::tempdir().unwrap();
let cwd_b = std::fs::canonicalize(cwd_b.path()).unwrap();
let mut w1 = TraceWriter::new(sessions.clone(), TraceExtras::default());
w1.on_event(&init_event("sess-old", &cwd_a));
let mut w2 = TraceWriter::new(
sessions.clone(),
TraceExtras {
kind: Some("subagent".to_string()),
..Default::default()
},
);
w2.on_event(&init_event("sess-sub", &cwd_a));
let mut w3 = TraceWriter::new(sessions.clone(), TraceExtras::default());
w3.on_event(&init_event("sess-b", &cwd_b));
let latest = find_latest_rollout(&sessions, &cwd_a).unwrap();
assert!(latest.to_string_lossy().contains("sess-old"), "{latest:?}");
let latest_b = find_latest_rollout(&sessions, &cwd_b).unwrap();
assert!(latest_b.to_string_lossy().contains("sess-b"));
let empty = tempfile::tempdir().unwrap();
assert!(find_latest_rollout(&sessions, empty.path()).is_none());
let by_id = find_rollout_by_id(&sessions, &cwd_b, "sess-sub").unwrap();
assert!(by_id.to_string_lossy().contains("sess-sub"));
assert!(find_rollout_by_id(&sessions, &cwd_b, "sess-nope").is_none());
}
#[test]
fn usage_records_round_trip_for_exact_resume() {
use locode_protocol::{Report, Status, Usage};
let root = tempfile::tempdir().unwrap();
let cwd = tempfile::tempdir().unwrap();
let cwd = std::fs::canonicalize(cwd.path()).unwrap();
let mut writer = TraceWriter::new(root.path().join("sessions"), TraceExtras::default());
writer.on_event(&init_event("sess-u", &cwd));
writer.on_event(&message_event(Role::User, "hi"));
let report = |input: u64, output: u64| Report {
schema_version: 1,
status: Status::Completed,
harness: "grok".into(),
api_schema: "mock".into(),
final_message: None,
structured_output: None,
turns: 1,
tool_calls: vec![],
usage: Usage {
input_tokens: input,
output_tokens: output,
cache_read_tokens: Some(7),
..Default::default()
},
session_id: "sess-u".into(),
stop_reason: None,
error: None,
};
writer.on_event(&locode_protocol::Event::Result {
report: report(100, 10),
});
writer.on_event(&locode_protocol::Event::Result {
report: report(200, 20),
});
assert!(writer.take_error().is_none());
let contents = read_rollout(writer.path().unwrap()).unwrap();
let usage = contents.last_usage.expect("usage recovered");
assert_eq!(usage.input_tokens, 200, "the last run's usage");
assert_eq!(usage.output_tokens, 20);
assert_eq!(usage.cache_read_tokens, Some(7));
let bare = root.path().join("bare.jsonl");
let meta_line = std::fs::read_to_string(writer.path().unwrap())
.unwrap()
.lines()
.next()
.unwrap()
.to_string();
std::fs::write(&bare, meta_line + "\n").unwrap();
assert!(read_rollout(&bare).unwrap().last_usage.is_none());
}
#[test]
fn reserved_kinds_and_parents_serialize() {
let root = tempfile::tempdir().unwrap();
let cwd = tempfile::tempdir().unwrap();
let cwd = std::fs::canonicalize(cwd.path()).unwrap();
let mut writer = TraceWriter::new(
root.path().join("sessions"),
TraceExtras {
cli_version: "x".to_string(),
kind: Some("subagent".to_string()),
parent_id: Some("sess-parent".to_string()),
group: Some("wf-1".to_string()),
..Default::default()
},
);
writer.on_event(&init_event("sess-4", &cwd));
let lines = read_lines(writer.path().unwrap());
let meta = &lines[0]["payload"];
assert_eq!(meta["kind"], "subagent");
assert_eq!(meta["parent_id"], "sess-parent");
assert_eq!(meta["group"], "wf-1");
let parsed: SessionMeta = serde_json::from_value(meta.clone()).unwrap();
assert_eq!(parsed.kind, "subagent");
}
}