use std::path::{Path, PathBuf};
use color_eyre::{eyre::bail, eyre::WrapErr, Result};
use serde::{Deserialize, Serialize};
use crate::agent::{ContentPart, Message};
pub const FORMAT_VERSION: u32 = 1;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum TurnEnd {
Complete,
Failed,
Interrupted,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "event", rename_all = "snake_case")]
pub enum SessionEvent {
TurnStart,
UserMessage {
text: String,
},
AssistantMessage {
blocks: Vec<ContentPart>,
},
ToolCall {
id: String,
name: String,
},
ToolResult {
id: String,
content: String,
#[serde(default)]
is_error: bool,
},
Compacted {
checkpoint: String,
replaced: usize,
},
TurnEnd {
reason: TurnEnd,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct Record {
seq: usize,
#[serde(flatten)]
event: SessionEvent,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionHeader {
pub kind: String,
pub version: u32,
pub id: String,
pub created_at: String,
pub cwd: Option<String>,
}
const TOOL_OUTCOME_UNKNOWN: &str =
"Procyon exited while this tool was running, so its outcome is unknown. Retry it only if it \
is read-only or idempotent; if it may have had an effect, verify the current state first or \
ask the user.";
const TOOL_NOT_STARTED: &str =
"Procyon exited before this tool started. Nothing happened, so retry it if it is still needed.";
pub fn interrupted_turn_closers(events: &[SessionEvent]) -> Vec<SessionEvent> {
let mut open_turn = false;
let mut pending: Vec<(String, bool)> = Vec::new();
for event in events {
match event {
SessionEvent::TurnStart => {
open_turn = true;
pending.clear();
}
SessionEvent::TurnEnd { .. } => {
open_turn = false;
pending.clear();
}
SessionEvent::AssistantMessage { blocks } => {
for block in blocks {
if let ContentPart::ToolUse { id, .. } = block {
pending.push((id.clone(), false));
}
}
}
SessionEvent::ToolCall { id, .. } => {
if let Some(entry) = pending.iter_mut().find(|(pid, _)| pid == id) {
entry.1 = true;
}
}
SessionEvent::ToolResult { id, .. } => {
pending.retain(|(pid, _)| pid != id);
}
_ => {}
}
}
if !open_turn {
return Vec::new();
}
let mut closers: Vec<SessionEvent> = pending
.into_iter()
.map(|(id, started)| SessionEvent::ToolResult {
id,
content: if started {
TOOL_OUTCOME_UNKNOWN.to_string()
} else {
TOOL_NOT_STARTED.to_string()
},
is_error: true,
})
.collect();
closers.push(SessionEvent::TurnEnd {
reason: TurnEnd::Interrupted,
});
closers
}
pub fn fold_to_messages(events: &[SessionEvent]) -> Vec<Message> {
let mut messages = Vec::new();
let mut pending_results: Vec<(String, String)> = Vec::new();
let flush = |messages: &mut Vec<Message>, results: &mut Vec<(String, String)>| {
if !results.is_empty() {
messages.push(Message::tool_results(std::mem::take(results)));
}
};
for event in events {
match event {
SessionEvent::ToolResult { id, content, .. } => {
pending_results.push((id.clone(), content.clone()));
continue;
}
_ => flush(&mut messages, &mut pending_results),
}
match event {
SessionEvent::UserMessage { text } => messages.push(Message::user(text)),
SessionEvent::AssistantMessage { blocks } => {
messages.push(Message::assistant(blocks.clone()))
}
SessionEvent::Compacted {
checkpoint,
replaced,
} => {
let cut = (*replaced).min(messages.len());
messages.drain(0..cut);
messages.insert(0, Message::user(checkpoint));
}
_ => {}
}
}
flush(&mut messages, &mut pending_results);
messages
}
fn slug(path: &Path) -> String {
let raw = path.to_string_lossy();
let mut out: String = raw
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect();
if out.len() > 120 {
out = out[out.len() - 120..].to_string();
}
let trimmed = out.trim_start_matches('-');
if trimmed.is_empty() {
"no-cwd".to_string()
} else {
trimmed.to_string()
}
}
pub fn sessions_root() -> Result<PathBuf> {
let base = dirs::data_dir()
.ok_or_else(|| color_eyre::eyre::eyre!("Failed to locate a data directory"))?;
Ok(base.join("procyon").join("sessions"))
}
fn session_dir(cwd: &Path) -> Result<PathBuf> {
Ok(sessions_root()?.join(slug(cwd)))
}
fn session_dir_under(root: &Path, cwd: &Path) -> PathBuf {
root.join(slug(cwd))
}
pub struct SessionLog {
path: PathBuf,
id: String,
next_seq: usize,
open_turn: bool,
file: Option<tokio::fs::File>,
}
impl SessionLog {
pub fn id(&self) -> &str {
&self.id
}
#[allow(dead_code)]
pub fn path(&self) -> &Path {
&self.path
}
pub async fn create(cwd: &Path) -> Result<Self> {
Self::create_under(&sessions_root()?, cwd).await
}
pub async fn create_under(root: &Path, cwd: &Path) -> Result<Self> {
let dir = session_dir_under(root, cwd);
tokio::fs::create_dir_all(&dir).await?;
let created_at = chrono::Utc::now();
let id = format!(
"{}-{}",
created_at.format("%Y%m%dT%H%M%S"),
std::process::id()
);
let path = dir.join(format!("{}.jsonl", id));
let header = SessionHeader {
kind: "session".to_string(),
version: FORMAT_VERSION,
id: id.clone(),
created_at: created_at.to_rfc3339(),
cwd: Some(cwd.to_string_lossy().to_string()),
};
let mut log = Self {
path,
id,
next_seq: 0,
open_turn: false,
file: None,
};
let mut line = serde_json::to_string(&header)?;
line.push('\n');
log.open_for_append().await?;
log.write_raw(&line).await?;
log.flush().await?;
Ok(log)
}
async fn open_for_append(&mut self) -> Result<()> {
let file = tokio::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&self.path)
.await
.wrap_err_with(|| format!("Failed to open {}", self.path.display()))?;
self.file = Some(file);
Ok(())
}
async fn write_raw(&mut self, line: &str) -> Result<()> {
use tokio::io::AsyncWriteExt;
if let Some(file) = self.file.as_mut() {
file.write_all(line.as_bytes()).await?;
}
Ok(())
}
pub async fn flush(&mut self) -> Result<()> {
use tokio::io::AsyncWriteExt;
if let Some(file) = self.file.as_mut() {
file.flush().await?;
file.sync_data().await?;
}
Ok(())
}
pub async fn append(&mut self, event: SessionEvent) -> Result<()> {
match &event {
SessionEvent::TurnStart => self.open_turn = true,
SessionEvent::TurnEnd { .. } => self.open_turn = false,
_ if !self.open_turn => {
bail!("{:?} appended outside a turn", event);
}
_ => {}
}
let record = Record {
seq: self.next_seq,
event,
};
self.next_seq += 1;
let mut line = serde_json::to_string(&record)?;
line.push('\n');
self.write_raw(&line).await
}
}
pub struct LoadedSession {
pub header: SessionHeader,
pub events: Vec<SessionEvent>,
pub repaired: Vec<SessionEvent>,
}
impl LoadedSession {
pub fn messages(&self) -> Vec<Message> {
let mut all = self.events.clone();
all.extend(self.repaired.clone());
fold_to_messages(&all)
}
pub fn transcript(&self) -> Vec<crate::channels::TranscriptEntry> {
use crate::channels::TranscriptEntry as Entry;
let mut out = Vec::new();
let mut names: std::collections::HashMap<&str, &str> = std::collections::HashMap::new();
for event in self.events.iter().chain(self.repaired.iter()) {
match event {
SessionEvent::UserMessage { text } => out.push(Entry::User(text.clone())),
SessionEvent::AssistantMessage { blocks } => {
for block in blocks {
if let ContentPart::Text { text } = block {
if !text.trim().is_empty() {
out.push(Entry::Agent(text.clone()));
}
}
}
}
SessionEvent::ToolCall { id, name } => {
names.insert(id.as_str(), name.as_str());
}
SessionEvent::ToolResult { id, is_error, .. } => out.push(Entry::Tool {
name: names
.get(id.as_str())
.copied()
.unwrap_or("tool")
.to_string(),
ok: !is_error,
}),
SessionEvent::Compacted { .. } => out.push(Entry::Compacted),
SessionEvent::TurnStart | SessionEvent::TurnEnd { .. } => {}
}
}
out
}
}
fn parse_log(content: &str) -> Result<(SessionHeader, Vec<SessionEvent>, usize)> {
let mut committed = 0usize;
let mut lines = content.split_inclusive('\n');
let header_line = lines
.next()
.ok_or_else(|| color_eyre::eyre::eyre!("Session log is empty"))?;
if !header_line.ends_with('\n') {
bail!("Session log has no complete header line");
}
let header: SessionHeader = serde_json::from_str(header_line.trim_end())
.wrap_err("Session log header is not readable")?;
if header.version != FORMAT_VERSION {
bail!(
"Session log is format version {}, but this build understands {}. Upgrade procyon.",
header.version,
FORMAT_VERSION
);
}
committed += header_line.len();
let mut events = Vec::new();
for line in lines {
if !line.ends_with('\n') {
break;
}
let trimmed = line.trim_end();
if trimmed.is_empty() {
committed += line.len();
continue;
}
let Ok(record) = serde_json::from_str::<Record>(trimmed) else {
break;
};
if record.seq != events.len() {
break;
}
events.push(record.event);
committed += line.len();
}
Ok((header, events, committed))
}
pub async fn load(path: &Path) -> Result<LoadedSession> {
let content = tokio::fs::read_to_string(path)
.await
.wrap_err_with(|| format!("Failed to read {}", path.display()))?;
let (header, events, committed) = parse_log(&content)?;
if committed < content.len() {
let file = tokio::fs::OpenOptions::new().write(true).open(path).await?;
file.set_len(committed as u64).await?;
file.sync_all().await?;
}
let repaired = interrupted_turn_closers(&events);
Ok(LoadedSession {
header,
events,
repaired,
})
}
pub struct Resumed {
pub log: SessionLog,
pub history: Vec<Message>,
pub transcript: Vec<crate::channels::TranscriptEntry>,
}
pub async fn resume(path: &Path) -> Result<Resumed> {
let loaded = load(path).await?;
let messages = loaded.messages();
let transcript = loaded.transcript();
let mut log = SessionLog {
path: path.to_path_buf(),
id: loaded.header.id.clone(),
next_seq: loaded.events.len(),
open_turn: !loaded.repaired.is_empty(),
file: None,
};
log.open_for_append().await?;
for event in loaded.repaired {
log.append(event).await?;
}
log.flush().await?;
Ok(Resumed {
log,
history: messages,
transcript,
})
}
pub async fn list(cwd: &Path) -> Result<Vec<(PathBuf, SessionHeader)>> {
match session_dir(cwd) {
Ok(dir) => list_under(&dir).await,
Err(_) => Ok(Vec::new()),
}
}
async fn list_under(dir: &Path) -> Result<Vec<(PathBuf, SessionHeader)>> {
if !tokio::fs::try_exists(dir).await.unwrap_or(false) {
return Ok(Vec::new());
}
let mut entries = tokio::fs::read_dir(&dir).await?;
let mut found = Vec::new();
while let Some(entry) = entries.next_entry().await? {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
continue;
}
let Ok(content) = tokio::fs::read_to_string(&path).await else {
continue;
};
let Some(first) = content.lines().next() else {
continue;
};
if let Ok(header) = serde_json::from_str::<SessionHeader>(first) {
found.push((path, header));
}
}
found.sort_by(|a, b| b.1.created_at.cmp(&a.1.created_at));
Ok(found)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn call(id: &str) -> ContentPart {
ContentPart::ToolUse {
id: id.to_string(),
name: "grep".to_string(),
input: json!({}),
}
}
#[test]
fn a_completed_turn_needs_no_repair() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "hi".to_string(),
},
SessionEvent::TurnEnd {
reason: TurnEnd::Complete,
},
];
assert!(interrupted_turn_closers(&events).is_empty());
}
#[test]
fn a_failed_turn_is_distinguishable_from_a_complete_one() {
let complete = serde_json::to_string(&SessionEvent::TurnEnd {
reason: TurnEnd::Complete,
})
.unwrap();
let failed = serde_json::to_string(&SessionEvent::TurnEnd {
reason: TurnEnd::Failed,
})
.unwrap();
assert!(complete.contains("\"complete\""), "got {}", complete);
assert!(failed.contains("\"failed\""), "got {}", failed);
assert_ne!(complete, failed);
}
#[test]
fn a_failed_turn_is_closed_and_needs_no_repair() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "deploy it".to_string(),
},
SessionEvent::TurnEnd {
reason: TurnEnd::Failed,
},
];
assert!(interrupted_turn_closers(&events).is_empty());
}
#[test]
fn an_open_turn_is_closed() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "hi".to_string(),
},
];
assert_eq!(
interrupted_turn_closers(&events),
vec![SessionEvent::TurnEnd {
reason: TurnEnd::Interrupted
}]
);
}
#[test]
fn a_call_that_started_is_marked_outcome_unknown() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::AssistantMessage {
blocks: vec![call("a")],
},
SessionEvent::ToolCall {
id: "a".to_string(),
name: "grep".to_string(),
},
];
let closers = interrupted_turn_closers(&events);
match &closers[0] {
SessionEvent::ToolResult {
id,
content,
is_error,
} => {
assert_eq!(id, "a");
assert!(*is_error);
assert!(content.contains("outcome is unknown"), "got {}", content);
assert!(content.contains("verify the current state"));
}
other => panic!("expected a tool result, got {:?}", other),
}
}
#[test]
fn a_call_that_never_started_is_safe_to_retry() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::AssistantMessage {
blocks: vec![call("a")],
},
];
let closers = interrupted_turn_closers(&events);
match &closers[0] {
SessionEvent::ToolResult { content, .. } => {
assert!(content.contains("Nothing happened"), "got {}", content);
}
other => panic!("expected a tool result, got {:?}", other),
}
}
#[test]
fn an_answered_call_is_not_closed_again() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::AssistantMessage {
blocks: vec![call("a")],
},
SessionEvent::ToolCall {
id: "a".to_string(),
name: "grep".to_string(),
},
SessionEvent::ToolResult {
id: "a".to_string(),
content: "ok".to_string(),
is_error: false,
},
];
assert_eq!(
interrupted_turn_closers(&events),
vec![SessionEvent::TurnEnd {
reason: TurnEnd::Interrupted
}]
);
}
#[test]
fn every_parallel_call_gets_its_own_closer() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::AssistantMessage {
blocks: vec![call("a"), call("b"), call("c")],
},
SessionEvent::ToolResult {
id: "b".to_string(),
content: "ok".to_string(),
is_error: false,
},
];
let closers = interrupted_turn_closers(&events);
assert_eq!(
closers.len(),
3,
"two results plus the turn end: {:?}",
closers
);
}
#[test]
fn a_repaired_transcript_pairs_every_call_with_a_result() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "do it".to_string(),
},
SessionEvent::AssistantMessage {
blocks: vec![call("a"), call("b")],
},
];
let mut all = events.clone();
all.extend(interrupted_turn_closers(&events));
let messages = fold_to_messages(&all);
let calls: usize = messages
.iter()
.flat_map(|m| &m.content)
.filter(|b| matches!(b, ContentPart::ToolUse { .. }))
.count();
let results: usize = messages
.iter()
.flat_map(|m| &m.content)
.filter(|b| matches!(b, ContentPart::ToolResult { .. }))
.count();
assert_eq!(calls, results, "the API rejects an unpaired transcript");
}
#[test]
fn results_of_one_turn_fold_into_a_single_user_message() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::AssistantMessage {
blocks: vec![call("a"), call("b")],
},
SessionEvent::ToolResult {
id: "a".to_string(),
content: "ra".to_string(),
is_error: false,
},
SessionEvent::ToolResult {
id: "b".to_string(),
content: "rb".to_string(),
is_error: false,
},
];
let messages = fold_to_messages(&events);
assert_eq!(messages.len(), 2, "{:?}", messages);
assert_eq!(messages[1].role, crate::agent::Role::User);
assert_eq!(messages[1].content.len(), 2);
}
fn text_of(message: &Message) -> &str {
match &message.content[0] {
ContentPart::Text { text } => text,
other => panic!("expected text, got {:?}", other),
}
}
#[test]
fn a_checkpoint_replaces_only_the_compacted_head() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "oldest".to_string(),
},
SessionEvent::UserMessage {
text: "older".to_string(),
},
SessionEvent::UserMessage {
text: "kept".to_string(),
},
SessionEvent::Compacted {
checkpoint: "CHECKPOINT".to_string(),
replaced: 2,
},
];
let messages = fold_to_messages(&events);
let rendered: Vec<&str> = messages.iter().map(text_of).collect();
assert_eq!(
rendered,
vec!["CHECKPOINT", "kept"],
"the retained tail must survive compaction"
);
}
#[test]
fn a_second_checkpoint_folds_over_the_first() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "a".to_string(),
},
SessionEvent::UserMessage {
text: "b".to_string(),
},
SessionEvent::Compacted {
checkpoint: "FIRST".to_string(),
replaced: 1,
},
SessionEvent::UserMessage {
text: "c".to_string(),
},
SessionEvent::Compacted {
checkpoint: "SECOND".to_string(),
replaced: 2,
},
];
let messages = fold_to_messages(&events);
let rendered: Vec<&str> = messages.iter().map(text_of).collect();
assert_eq!(rendered, vec!["SECOND", "c"]);
}
#[test]
fn a_checkpoint_claiming_more_than_exists_does_not_panic() {
let events = vec![
SessionEvent::TurnStart,
SessionEvent::UserMessage {
text: "only".to_string(),
},
SessionEvent::Compacted {
checkpoint: "CHECKPOINT".to_string(),
replaced: 99,
},
];
let messages = fold_to_messages(&events);
assert_eq!(messages.len(), 1);
assert_eq!(text_of(&messages[0]), "CHECKPOINT");
}
fn header_line() -> String {
serde_json::to_string(&SessionHeader {
kind: "session".to_string(),
version: FORMAT_VERSION,
id: "s1".to_string(),
created_at: "2026-08-21T00:00:00Z".to_string(),
cwd: None,
})
.unwrap()
}
#[test]
fn a_torn_final_line_is_discarded() {
let good = serde_json::to_string(&Record {
seq: 0,
event: SessionEvent::TurnStart,
})
.unwrap();
let content = format!(
"{}\n{}\n{{\"seq\":1,\"event\":\"user_mes",
header_line(),
good
);
let (_, events, committed) = parse_log(&content).unwrap();
assert_eq!(events.len(), 1, "the torn record must not be parsed");
assert!(committed < content.len(), "the tail must be excluded");
}
#[test]
fn a_gap_in_the_sequence_stops_the_read() {
let first = serde_json::to_string(&Record {
seq: 0,
event: SessionEvent::TurnStart,
})
.unwrap();
let skipped = serde_json::to_string(&Record {
seq: 5,
event: SessionEvent::TurnStart,
})
.unwrap();
let content = format!("{}\n{}\n{}\n", header_line(), first, skipped);
let (_, events, _) = parse_log(&content).unwrap();
assert_eq!(events.len(), 1);
}
#[test]
fn a_newer_format_is_refused_with_an_upgrade_message() {
let mut header: serde_json::Value = serde_json::from_str(&header_line()).unwrap();
header["version"] = json!(FORMAT_VERSION + 1);
let content = format!("{}\n", header);
let err = parse_log(&content).unwrap_err().to_string();
assert!(err.contains("Upgrade procyon"), "got {}", err);
}
#[tokio::test]
async fn appending_outside_a_turn_is_refused() {
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog {
path: temp.path().join("s.jsonl"),
id: "s".to_string(),
next_seq: 0,
open_turn: false,
file: None,
};
log.open_for_append().await.unwrap();
let err = log
.append(SessionEvent::UserMessage {
text: "stray".to_string(),
})
.await
.unwrap_err()
.to_string();
assert!(err.contains("outside a turn"), "got {}", err);
}
#[tokio::test]
async fn a_written_log_round_trips() {
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog::create_under(temp.path(), temp.path())
.await
.unwrap();
log.append(SessionEvent::TurnStart).await.unwrap();
log.append(SessionEvent::UserMessage {
text: "olá".to_string(),
})
.await
.unwrap();
log.append(SessionEvent::TurnEnd {
reason: TurnEnd::Complete,
})
.await
.unwrap();
log.flush().await.unwrap();
let loaded = load(log.path()).await.unwrap();
assert_eq!(loaded.events.len(), 3);
assert!(loaded.repaired.is_empty());
let messages = loaded.messages();
assert_eq!(messages.len(), 1);
assert!(matches!(&messages[0].content[0], ContentPart::Text { text } if text == "olá"));
}
#[tokio::test]
async fn resuming_a_crashed_log_closes_the_turn_on_disk() {
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog::create_under(temp.path(), temp.path())
.await
.unwrap();
log.append(SessionEvent::TurnStart).await.unwrap();
log.append(SessionEvent::UserMessage {
text: "do it".to_string(),
})
.await
.unwrap();
log.append(SessionEvent::AssistantMessage {
blocks: vec![call("a")],
})
.await
.unwrap();
log.flush().await.unwrap();
let path = log.path().to_path_buf();
drop(log);
let messages = resume(&path).await.unwrap().history;
let reloaded = load(&path).await.unwrap();
assert!(
reloaded.repaired.is_empty(),
"a resumed log must no longer need repair"
);
assert!(matches!(
reloaded.events.last(),
Some(SessionEvent::TurnEnd {
reason: TurnEnd::Interrupted
})
));
let results: usize = messages
.iter()
.flat_map(|m| &m.content)
.filter(|b| matches!(b, ContentPart::ToolResult { .. }))
.count();
assert_eq!(results, 1, "the dead call must be answered");
}
#[tokio::test]
async fn a_log_truncated_mid_write_is_recovered_and_resumable() {
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog::create_under(temp.path(), temp.path())
.await
.unwrap();
let path = log.path().to_path_buf();
log.append(SessionEvent::TurnStart).await.unwrap();
log.append(SessionEvent::UserMessage {
text: "deploy it".to_string(),
})
.await
.unwrap();
log.append(SessionEvent::AssistantMessage {
blocks: vec![call("a")],
})
.await
.unwrap();
log.append(SessionEvent::ToolCall {
id: "a".to_string(),
name: "caatinga_deploy".to_string(),
})
.await
.unwrap();
log.flush().await.unwrap();
drop(log);
let mut raw = tokio::fs::read_to_string(&path).await.unwrap();
raw.push_str("{\"seq\":4,\"event\":\"tool_res");
tokio::fs::write(&path, &raw).await.unwrap();
let messages = resume(&path).await.unwrap().history;
let on_disk = tokio::fs::read_to_string(&path).await.unwrap();
assert!(
!on_disk.contains("tool_res\""),
"the torn record survived: {}",
on_disk
);
let flagged = messages
.iter()
.flat_map(|m| &m.content)
.filter_map(|b| match b {
ContentPart::ToolResult { content, .. } => Some(content.as_str()),
_ => None,
})
.any(|c| c.contains("outcome is unknown"));
assert!(
flagged,
"a started tool must not look untouched: {:?}",
messages
);
let reloaded = load(&path).await.unwrap();
assert!(
reloaded.repaired.is_empty(),
"resume must leave a clean log"
);
}
#[tokio::test]
async fn a_resumed_transcript_carries_both_sides_of_the_conversation() {
use crate::channels::TranscriptEntry as Entry;
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(temp.path()).await.unwrap();
log.append(SessionEvent::TurnStart).await.unwrap();
log.append(SessionEvent::UserMessage {
text: "liste os arquivos".to_string(),
})
.await
.unwrap();
log.append(SessionEvent::ToolCall {
id: "t1".to_string(),
name: "list_dir".to_string(),
})
.await
.unwrap();
log.append(SessionEvent::ToolResult {
id: "t1".to_string(),
content: "a.rs".to_string(),
is_error: false,
})
.await
.unwrap();
log.append(SessionEvent::AssistantMessage {
blocks: vec![ContentPart::Text {
text: "sao estes".to_string(),
}],
})
.await
.unwrap();
log.append(SessionEvent::TurnEnd {
reason: TurnEnd::Complete,
})
.await
.unwrap();
log.flush().await.unwrap();
let path = log.path().to_path_buf();
drop(log);
let transcript = resume(&path).await.unwrap().transcript;
assert_eq!(
transcript,
vec![
Entry::User("liste os arquivos".to_string()),
Entry::Tool {
name: "list_dir".to_string(),
ok: true
},
Entry::Agent("sao estes".to_string()),
]
);
}
#[tokio::test]
async fn a_resumed_transcript_remembers_which_calls_failed() {
use crate::channels::TranscriptEntry as Entry;
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog::create(temp.path()).await.unwrap();
log.append(SessionEvent::TurnStart).await.unwrap();
log.append(SessionEvent::ToolCall {
id: "t1".to_string(),
name: "caatinga_deploy".to_string(),
})
.await
.unwrap();
log.append(SessionEvent::ToolResult {
id: "t1".to_string(),
content: "Error: boom".to_string(),
is_error: true,
})
.await
.unwrap();
log.flush().await.unwrap();
let path = log.path().to_path_buf();
drop(log);
let transcript = resume(&path).await.unwrap().transcript;
assert!(transcript.contains(&Entry::Tool {
name: "caatinga_deploy".to_string(),
ok: false
}));
}
#[tokio::test]
async fn a_resumed_compacted_session_matches_the_live_history() {
let temp = tempfile::tempdir().unwrap();
let mut log = SessionLog::create_under(temp.path(), temp.path())
.await
.unwrap();
let path = log.path().to_path_buf();
let mut live: Vec<Message> = Vec::new();
log.append(SessionEvent::TurnStart).await.unwrap();
for i in 0..3 {
let user = format!("question {}", i);
let reply = format!("answer {}", i);
log.append(SessionEvent::UserMessage { text: user.clone() })
.await
.unwrap();
live.push(Message::user(&user));
log.append(SessionEvent::AssistantMessage {
blocks: vec![ContentPart::Text {
text: reply.clone(),
}],
})
.await
.unwrap();
live.push(Message::assistant(vec![ContentPart::Text { text: reply }]));
}
let cut = 4;
let checkpoint = "CHECKPOINT: earlier turns summarized".to_string();
log.append(SessionEvent::Compacted {
checkpoint: checkpoint.clone(),
replaced: cut,
})
.await
.unwrap();
live.splice(0..cut, std::iter::once(Message::user(&checkpoint)));
log.append(SessionEvent::TurnEnd {
reason: TurnEnd::Complete,
})
.await
.unwrap();
log.flush().await.unwrap();
drop(log);
let restored = resume(&path).await.unwrap().history;
let rendered = |messages: &[Message]| -> Vec<(String, String)> {
messages
.iter()
.map(|m| (m.role.to_string(), text_of(m).to_string()))
.collect()
};
assert_eq!(
rendered(&restored),
rendered(&live),
"a resumed compacted session must match the live history"
);
assert_eq!(restored.len(), 3, "checkpoint plus the retained tail");
}
#[tokio::test]
async fn listing_reports_newest_first_and_ignores_other_files() {
let temp = tempfile::tempdir().unwrap();
let first = SessionLog::create_under(temp.path(), temp.path())
.await
.unwrap();
tokio::fs::write(first.path().parent().unwrap().join("notes.txt"), "x")
.await
.unwrap();
let listed = list_under(&session_dir_under(temp.path(), temp.path()))
.await
.unwrap();
assert_eq!(listed.len(), 1, "only .jsonl logs count: {:?}", listed);
assert_eq!(listed[0].1.id, first.id());
}
#[tokio::test]
async fn a_log_stays_inside_the_root_it_was_given() {
let temp = tempfile::tempdir().unwrap();
let log = SessionLog::create_under(temp.path(), temp.path())
.await
.unwrap();
assert!(
log.path().starts_with(temp.path()),
"{} escaped {}",
log.path().display(),
temp.path().display()
);
assert!(
!log.path().starts_with(sessions_root().unwrap()),
"a test log must not land in the real sessions directory"
);
}
#[test]
fn the_slug_keeps_paths_filesystem_safe() {
let s = slug(Path::new("/home/user/My Project (v2)"));
assert!(
!s.contains('/') && !s.contains(' ') && !s.contains('('),
"got {}",
s
);
}
#[test]
fn the_slug_does_not_start_with_a_dash() {
for path in ["/home/user/proj", "/tmp/.tmpXYZ", "///weird"] {
let s = slug(Path::new(path));
assert!(!s.starts_with('-'), "{} produced {}", path, s);
}
}
#[test]
fn an_unnameable_path_still_yields_a_directory() {
assert_eq!(slug(Path::new("/")), "no-cwd");
}
}