use std::collections::{HashMap, HashSet};
use std::fs::{File, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use crate::chat::turns::{ChatRole, ChatTurn};
use crate::chat::types::{ConversationItem, Lifecycle};
use crate::project_session::ProjectObservation;
use crate::task::TaskObservation;
use crate::wave::playhead::{BodyProvenance, Playhead, PlayheadEvent};
use crate::wave::state::LoopState;
const FORMAT_VERSION: u32 = 1;
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct MessageId(pub String);
impl std::fmt::Display for MessageId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MessageOp {
Message,
Steer,
Interrupt,
Say,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Usage {
pub input_tokens: Option<u64>,
pub output_tokens: Option<u64>,
pub cache_read_tokens: Option<u64>,
pub cost_usd: Option<f64>,
}
impl Usage {
pub fn empty() -> Self {
Self {
input_tokens: None,
output_tokens: None,
cache_read_tokens: None,
cost_usd: None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WorkerOutcome {
Completed,
Failed,
}
impl WorkerOutcome {
pub fn name(&self) -> &'static str {
match self {
Self::Completed => "completed",
Self::Failed => "failed",
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct PendingMessage {
pub id: MessageId,
pub op: MessageOp,
pub text: String,
pub from: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Event {
pub v: u32,
pub seq: u64,
#[serde(with = "time::serde::rfc3339")]
pub at: OffsetDateTime,
pub kind: EventKind,
}
impl Event {
pub fn at_rfc3339(&self) -> String {
self.at
.format(&time::format_description::well_known::Rfc3339)
.expect("journal timestamps are representable as RFC 3339")
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum EventKind {
UserMessage {
id: MessageId,
op: MessageOp,
text: String,
from: Option<String>,
},
TurnStarted {
turn_id: String,
answers: Vec<MessageId>,
body: Option<Box<BodyProvenance>>,
},
TurnItem {
turn_id: String,
item: ConversationItem,
},
TurnSteered {
turn_id: String,
answers: Vec<MessageId>,
},
MessagesRequeued {
ids: Vec<MessageId>,
},
TurnFinished {
turn_id: String,
status: Lifecycle,
usage: Usage,
termination_reason: Option<String>,
},
LoopState {
from: LoopState,
to: LoopState,
reason: String,
},
PlayheadChanged {
event: PlayheadEvent,
playhead: Box<Playhead>,
},
RunObserved {
run_id: String,
session_id: String,
flow: String,
task: String,
},
RunCompleted {
run_id: String,
outcome: WorkerOutcome,
summary: String,
},
TaskObserved {
observation: TaskObservation,
},
ProjectObserved {
observation: ProjectObservation,
},
ChannelOpened {
name: String,
run_id: String,
},
MemoryUpdated {
summary: String,
},
MemoryAdded {
fact: String,
},
ServerStarted {
pid: u32,
endpoint: String,
},
}
pub fn journal_path(repo_root: &Path, wave: &str) -> PathBuf {
repo_root
.join(".lf")
.join("journal")
.join("waves")
.join(wave)
.join("journal.jsonl")
}
pub fn read_events(path: &Path) -> Vec<Event> {
let Ok(raw) = std::fs::read_to_string(path) else {
return Vec::new();
};
let mut events = Vec::new();
for line in raw.lines() {
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
match serde_json::from_str::<Event>(trimmed) {
Ok(event) if event.v == FORMAT_VERSION => events.push(event),
Ok(event) => {
tracing::warn!(
path = %path.display(),
version = event.v,
"journal line from another format version; stopping read-only fold"
);
break;
}
Err(_) => break,
}
}
events
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum NarrationLevel {
Info,
Debug,
}
#[derive(Debug)]
struct Narration {
level: NarrationLevel,
line: String,
}
fn info(line: String) -> Narration {
Narration {
level: NarrationLevel::Info,
line,
}
}
fn debug(line: String) -> Narration {
Narration {
level: NarrationLevel::Debug,
line,
}
}
#[derive(Debug)]
struct TurnNarration {
turn_id: String,
items: usize,
text_shown: bool,
}
#[derive(Debug, Default)]
struct Narrator {
turns: Vec<TurnNarration>,
}
impl Narrator {
fn narrate(&mut self, kind: &EventKind) {
let narration = self.render(kind);
match narration.level {
NarrationLevel::Info => tracing::info!("{}", narration.line),
NarrationLevel::Debug => tracing::debug!("{}", narration.line),
}
}
fn render(&mut self, kind: &EventKind) -> Narration {
match kind {
EventKind::UserMessage { id, op, text, from } => {
let op_tag = match op {
MessageOp::Message | MessageOp::Say => "",
MessageOp::Steer => "(steer) ",
MessageOp::Interrupt => "(interrupt) ",
};
let byline = from
.as_ref()
.map(|from| format!("[{from}] "))
.unwrap_or_default();
info(format!(
"chat ← {op_tag}{byline}\"{}\" ({id})",
ellipsize(text, 60)
))
}
EventKind::TurnStarted {
turn_id, answers, ..
} => {
let turn = self.turn_mut(turn_id);
turn.items = 0;
turn.text_shown = false;
info(format!("turn {turn_id} opened{}", answers_segment(answers)))
}
EventKind::TurnItem { turn_id, item } => {
let turn = self.turn_mut(turn_id);
turn.items += 1;
match item {
ConversationItem::Command {
command, status, ..
} => info(format!(
" $ {} → {}",
ellipsize(&command.join(" "), 70),
status.name()
)),
ConversationItem::Message { text, .. } => {
if turn.text_shown {
debug(format!(" loop: \"{}\"", ellipsize(text, 120)))
} else {
turn.text_shown = true;
info(format!("loop: \"{}\"", ellipsize(text, 80)))
}
}
ConversationItem::Thought { text, .. } => {
debug(format!(" thought: \"{}\"", ellipsize(text, 120)))
}
ConversationItem::File {
changes, status, ..
} => {
let what = match changes.as_slice() {
[only] => only.path.clone(),
many => format!("{} files", many.len()),
};
info(format!(" edit {what} → {}", status.name()))
}
ConversationItem::Tool { name, status, .. } => {
info(format!(" tool {name} → {}", status.name()))
}
}
}
EventKind::TurnSteered { turn_id, answers } => info(format!(
"turn {turn_id} steered{}",
answers_segment(answers)
)),
EventKind::MessagesRequeued { ids } => {
let ids: Vec<&str> = ids.iter().map(|id| id.0.as_str()).collect();
info(format!("messages requeued: {}", ids.join(", ")))
}
EventKind::TurnFinished {
turn_id,
status,
usage,
..
} => {
let items = self.finish_turn(turn_id);
let plural = if items == 1 { "" } else { "s" };
info(format!(
"turn {turn_id} {} · {items} item{plural}{}",
status.name(),
usage_segment(usage)
))
}
EventKind::LoopState { from, to, reason } => {
info(format!("state {} → {} ({reason})", from.name(), to.name()))
}
EventKind::PlayheadChanged { event, .. } => match event {
PlayheadEvent::FlowEnqueued { flow, .. } => {
info(format!("playhead enqueued · {flow}"))
}
PlayheadEvent::InvocationStarted { flow, .. } => {
info(format!("playhead entered · {flow}"))
}
PlayheadEvent::InvocationCompleted { flow, .. } => {
info(format!("playhead completed · {flow}"))
}
PlayheadEvent::StepStarted { step, .. } => {
info(format!("playhead now · {} / {}", step.flow, step.step))
}
PlayheadEvent::BodySessionUpdated { session_id, .. } => {
info(format!("playhead session · {session_id}"))
}
PlayheadEvent::StepFinished {
step,
outcome,
reason,
..
} => info(format!(
"playhead {} · {} / {} ({reason})",
outcome.name(),
step.flow,
step.step
)),
},
EventKind::RunObserved {
run_id, flow, task, ..
} => info(format!(
"observed worker {} flow={flow} started · {}",
short_id(run_id),
ellipsize(task, 60)
)),
EventKind::RunCompleted {
run_id,
outcome,
summary,
} => info(format!(
"observed run {} {} · {}",
short_id(run_id),
outcome.name(),
ellipsize(summary, 60)
)),
EventKind::TaskObserved { observation } => info(format!(
"observed task {} event {} · {}",
observation.issue_identifier,
observation.event_id,
ellipsize(&observation.prompt(), 70)
)),
EventKind::ProjectObserved { observation } => info(format!(
"observed project {} event {} · {}",
observation.project,
observation.event_id,
ellipsize(&observation.prompt(), 70)
)),
EventKind::ChannelOpened { name, run_id } => {
info(format!("channel {name} opened · run {}", short_id(run_id)))
}
EventKind::MemoryUpdated { summary } => {
info(format!("memory curated: {}", ellipsize(summary, 70)))
}
EventKind::MemoryAdded { fact } => {
info(format!("memory added: {}", ellipsize(fact, 70)))
}
EventKind::ServerStarted { pid, endpoint } => {
info(format!("server started · pid {pid} · {endpoint}"))
}
}
}
fn turn_mut(&mut self, turn_id: &str) -> &mut TurnNarration {
if let Some(pos) = self.turns.iter().position(|t| t.turn_id == turn_id) {
return &mut self.turns[pos];
}
self.turns.push(TurnNarration {
turn_id: turn_id.to_string(),
items: 0,
text_shown: false,
});
self.turns.last_mut().expect("just pushed")
}
fn finish_turn(&mut self, turn_id: &str) -> usize {
match self.turns.iter().position(|t| t.turn_id == turn_id) {
Some(pos) => self.turns.remove(pos).items,
None => 0,
}
}
}
pub(crate) fn ellipsize(text: &str, max: usize) -> String {
let flat = text.split_whitespace().collect::<Vec<_>>().join(" ");
if flat.chars().count() <= max {
return flat;
}
let mut cut: String = flat.chars().take(max).collect();
cut.push('…');
cut
}
fn fmt_tokens(n: u64) -> String {
if n < 1000 {
return n.to_string();
}
let k = n as f64 / 1000.0;
if k < 10.0 {
format!("{k:.1}k")
} else {
format!("{k:.0}k")
}
}
fn answers_segment(answers: &[MessageId]) -> String {
if answers.is_empty() {
return String::new();
}
let ids: Vec<&str> = answers.iter().map(|id| id.0.as_str()).collect();
format!(" (answers: {})", ids.join(", "))
}
fn usage_segment(usage: &Usage) -> String {
let mut parts = Vec::new();
if let Some(input) = usage.input_tokens {
parts.push(format!("{} in", fmt_tokens(input)));
}
if let Some(output) = usage.output_tokens {
parts.push(format!("{} out", fmt_tokens(output)));
}
if parts.is_empty() {
return String::new();
}
let mut segment = format!(" · {}", parts.join(" / "));
if let Some(cached) = usage.cache_read_tokens {
segment.push_str(&format!(" ({} cached)", fmt_tokens(cached)));
}
segment
}
pub(crate) fn short_id(run_id: &str) -> String {
run_id.chars().take(8).collect()
}
#[derive(Debug)]
pub struct Journal {
file: File,
next_seq: u64,
narrator: Narrator,
}
impl Journal {
pub fn open(path: &Path) -> anyhow::Result<(Self, Vec<Event>)> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let raw = match std::fs::read_to_string(path) {
Ok(s) => s,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => String::new(),
Err(err) => return Err(err.into()),
};
let mut events = Vec::new();
let mut good_bytes = 0usize;
let mut offset = 0usize;
for line in raw.split_inclusive('\n') {
let start = offset;
offset += line.len();
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
match serde_json::from_str::<Event>(trimmed) {
Ok(event) if event.v == FORMAT_VERSION => {
events.push(event);
good_bytes = offset;
}
Ok(event) => {
anyhow::bail!(
"journal {} has format v{} at seq {}; this build reads v{FORMAT_VERSION}",
path.display(),
event.v,
event.seq,
);
}
Err(err) => {
tracing::warn!(
path = %path.display(),
byte_offset = start,
error = %err,
"journal has an unparseable tail (crash mid-write?); truncating to last parseable line"
);
break;
}
}
}
if good_bytes < raw.len() {
let file = OpenOptions::new().write(true).open(path)?;
file.set_len(good_bytes as u64)?;
}
let next_seq = events.last().map(|e| e.seq + 1).unwrap_or(1);
let file = OpenOptions::new().create(true).append(true).open(path)?;
Ok((
Self {
file,
next_seq,
narrator: Narrator::default(),
},
events,
))
}
#[cfg(test)]
pub fn next_seq(&self) -> u64 {
self.next_seq
}
pub fn append(&mut self, build: impl FnOnce(u64) -> EventKind) -> Event {
let event = Event {
v: FORMAT_VERSION,
seq: self.next_seq,
at: OffsetDateTime::now_utc(),
kind: build(self.next_seq),
};
self.next_seq += 1;
match serde_json::to_string(&event) {
Ok(line) => {
if let Err(err) = writeln!(self.file, "{line}").and_then(|_| self.file.flush()) {
tracing::error!(seq = event.seq, error = %err, "failed to append journal event");
} else {
self.narrator.narrate(&event.kind);
}
}
Err(err) => {
tracing::error!(seq = event.seq, error = %err, "failed to serialize journal event");
}
}
event
}
}
#[derive(Debug)]
pub struct ThreadFold {
pub turns: Vec<ChatTurn>,
pub open: Vec<ChatTurn>,
pub state: LoopState,
pub playhead: Option<Playhead>,
pub pending_messages: Vec<PendingMessage>,
pub messages: HashMap<MessageId, PendingMessage>,
pub tasks: HashMap<MessageId, TaskObservation>,
pub projects: HashMap<MessageId, ProjectObservation>,
pub open_claims: Vec<MessageId>,
pub memory_adds: Vec<String>,
}
fn legacy_channel_opened_turn(event: &Event, name: &str) -> ChatTurn {
let mut turn = ChatTurn::user(
format!("turn-{}", event.seq),
format!("work line {name} opened"),
);
turn.created_at = event.at_rfc3339();
turn.from = Some("worker".to_string());
turn
}
pub fn run_completed_turn(
event: &Event,
run_id: &str,
outcome: WorkerOutcome,
summary: &str,
) -> ChatTurn {
let mut text = format!("run {} {}", short_id(run_id), outcome.name());
if !summary.trim().is_empty() {
text.push_str(&format!(" · {}", summary.trim()));
}
let mut turn = ChatTurn::user(format!("turn-{}", event.seq), text);
turn.created_at = event.at_rfc3339();
turn.from = Some("observer".to_string());
turn
}
pub fn restore_pending(
pending: &mut Vec<PendingMessage>,
messages: &HashMap<MessageId, PendingMessage>,
ids: &[MessageId],
) -> Vec<MessageId> {
let mut restored = Vec::new();
for id in ids {
if pending.iter().any(|message| &message.id == id) {
continue;
}
let Some(message) = messages.get(id) else {
tracing::warn!(id = %id, "requeue of an unknown message id; dropped");
continue;
};
pending.push(message.clone());
restored.push(id.clone());
}
restored
}
pub fn fold_thread(events: &[Event]) -> ThreadFold {
let mut turns: Vec<ChatTurn> = Vec::new();
let mut open: Vec<ChatTurn> = Vec::new();
let mut state = LoopState::Idle;
let mut playhead: Option<Playhead> = None;
let mut pending_messages: Vec<PendingMessage> = Vec::new();
let mut messages: HashMap<MessageId, PendingMessage> = HashMap::new();
let mut tasks: HashMap<MessageId, TaskObservation> = HashMap::new();
let mut projects: HashMap<MessageId, ProjectObservation> = HashMap::new();
let mut consumed_messages: HashSet<MessageId> = HashSet::new();
let mut claims_by_open_turn: HashMap<String, Vec<MessageId>> = HashMap::new();
let mut memory_adds: Vec<String> = Vec::new();
for event in events {
match &event.kind {
EventKind::UserMessage { id, op, text, from } => {
let mut turn = ChatTurn::user(format!("turn-{}", event.seq), text.clone());
turn.created_at = event.at_rfc3339();
turn.from = from.clone();
turns.push(turn);
let message = PendingMessage {
id: id.clone(),
op: *op,
text: text.clone(),
from: from.clone(),
};
if !consumed_messages.contains(id) {
pending_messages.push(message.clone());
}
messages.insert(id.clone(), message);
}
EventKind::TaskObserved { observation } => {
let message = task_observation_message(observation);
let turn = ChatTurn::child_activity(
format!("turn-{}", event.seq),
event.at_rfc3339(),
"task".to_string(),
crate::chat::turns::ChildControlActivity::from_task(observation),
);
turns.push(turn);
if !consumed_messages.contains(&message.id) {
pending_messages.push(message.clone());
}
tasks.insert(message.id.clone(), observation.clone());
messages.insert(message.id.clone(), message);
}
EventKind::ProjectObserved { observation } => {
let message = project_observation_message(observation);
let turn = ChatTurn::child_activity(
format!("turn-{}", event.seq),
event.at_rfc3339(),
"project".to_string(),
crate::chat::turns::ChildControlActivity::from_project(observation),
);
turns.push(turn);
if !consumed_messages.contains(&message.id) {
pending_messages.push(message.clone());
}
projects.insert(message.id.clone(), observation.clone());
messages.insert(message.id.clone(), message);
}
EventKind::TurnStarted {
turn_id,
answers,
body,
} => {
mark_consumed(&mut pending_messages, &mut consumed_messages, answers);
claims_by_open_turn.insert(turn_id.clone(), answers.clone());
open.push(ChatTurn {
id: turn_id.clone(),
role: ChatRole::Assistant,
text: String::new(),
status: Lifecycle::Running,
items: Vec::new(),
created_at: event.at_rfc3339(),
from: None,
body: body.as_deref().cloned(),
activity: None,
});
}
EventKind::TurnItem { turn_id, item } => {
let Some(turn) = open.iter_mut().find(|t| &t.id == turn_id) else {
tracing::warn!(
turn_id,
seq = event.seq,
"TurnItem for a turn that isn't open"
);
continue;
};
turn.absorb_item(item.clone());
}
EventKind::TurnFinished {
turn_id,
status,
termination_reason,
..
} => {
let Some(pos) = open.iter().position(|t| &t.id == turn_id) else {
tracing::warn!(
turn_id,
seq = event.seq,
"TurnFinished for a turn that isn't open"
);
continue;
};
let mut turn = open.remove(pos);
turn.status = *status;
turn.close_body(event.at_rfc3339(), termination_reason.clone());
claims_by_open_turn.remove(turn_id);
turns.push(turn);
}
EventKind::LoopState { to, .. } => {
state = to.clone();
}
EventKind::PlayheadChanged {
event,
playhead: snapshot,
} => {
if let PlayheadEvent::BodySessionUpdated {
body_id,
session_id,
} = event
{
if let Some(body) = open
.iter_mut()
.filter_map(|turn| turn.body.as_mut())
.find(|body| &body.body_id == body_id)
{
body.session_id = Some(session_id.clone());
}
}
playhead = Some(snapshot.as_ref().clone());
}
EventKind::TurnSteered { turn_id, answers } => {
mark_consumed(&mut pending_messages, &mut consumed_messages, answers);
if let Some(claims) = claims_by_open_turn.get_mut(turn_id) {
claims.extend(answers.iter().cloned());
}
}
EventKind::MessagesRequeued { ids } => {
restore_pending(&mut pending_messages, &messages, ids);
}
EventKind::ChannelOpened { name, .. } => {
turns.push(legacy_channel_opened_turn(event, name));
}
EventKind::RunCompleted {
run_id,
outcome,
summary,
} => {
turns.push(run_completed_turn(event, run_id, *outcome, summary));
}
EventKind::MemoryAdded { fact } => {
memory_adds.push(fact.clone());
}
EventKind::MemoryUpdated { .. } => {
memory_adds.clear();
}
EventKind::RunObserved { .. } | EventKind::ServerStarted { .. } => {}
}
}
let open_claims = open
.iter()
.filter_map(|turn| claims_by_open_turn.remove(&turn.id))
.flatten()
.collect();
ThreadFold {
turns,
open,
state,
playhead,
pending_messages,
messages,
tasks,
projects,
open_claims,
memory_adds,
}
}
pub fn task_observation_message(observation: &TaskObservation) -> PendingMessage {
PendingMessage {
id: MessageId(observation.inbox_id()),
op: MessageOp::Message,
text: observation.prompt(),
from: Some("task".to_string()),
}
}
pub fn project_observation_message(observation: &ProjectObservation) -> PendingMessage {
PendingMessage {
id: MessageId(observation.inbox_id()),
op: MessageOp::Message,
text: observation.prompt(),
from: Some("project".to_string()),
}
}
fn mark_consumed(
pending_messages: &mut Vec<PendingMessage>,
consumed_messages: &mut HashSet<MessageId>,
answers: &[MessageId],
) {
if answers.is_empty() {
return;
}
for answer in answers {
consumed_messages.insert(answer.clone());
}
pending_messages.retain(|message| !consumed_messages.contains(&message.id));
}
#[cfg(test)]
mod tests {
use super::*;
fn open_tmp() -> (tempfile::TempDir, PathBuf) {
let tmp = tempfile::tempdir().expect("tempdir");
let path = journal_path(tmp.path(), "ship");
(tmp, path)
}
fn user_message(seq: u64, text: &str) -> EventKind {
EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op: MessageOp::Message,
text: text.to_string(),
from: None,
}
}
fn task_observation() -> crate::task::TaskObservation {
crate::task::TaskObservation {
session_id: crate::task::TaskSessionId::from_raw("ts_example"),
issue_identifier: "INF-123".to_string(),
event_id: 42,
control_source: None,
event: crate::task::TaskEventKind::Failed {
error: "provider stopped".to_string(),
resumable: true,
},
}
}
#[test]
fn append_then_open_replays_events_and_continues_seq() {
let (_tmp, path) = open_tmp();
{
let (mut journal, events) = Journal::open(&path).expect("open");
assert!(events.is_empty());
assert_eq!(journal.next_seq(), 1);
journal.append(|seq| user_message(seq, "hello"));
journal.append(|seq| user_message(seq, "again"));
}
let (journal, events) = Journal::open(&path).expect("reopen");
assert_eq!(events.len(), 2);
assert_eq!(events[0].seq, 1);
assert_eq!(events[1].seq, 2);
assert_eq!(journal.next_seq(), 3);
}
#[test]
fn corrupt_trailing_line_is_truncated_not_fatal() {
let (_tmp, path) = open_tmp();
{
let (mut journal, _) = Journal::open(&path).expect("open");
journal.append(|seq| user_message(seq, "kept"));
}
let mut raw = std::fs::read_to_string(&path).expect("read");
raw.push_str(r#"{"v":1,"seq":2,"at":"2026-"#);
std::fs::write(&path, &raw).expect("corrupt");
let (mut journal, events) = Journal::open(&path).expect("reopen tolerates tail");
assert_eq!(events.len(), 1);
assert_eq!(journal.next_seq(), 2);
journal.append(|seq| user_message(seq, "after crash"));
drop(journal);
let (_, events) = Journal::open(&path).expect("reopen again");
assert_eq!(events.len(), 2);
assert_eq!(events[1].seq, 2);
}
#[test]
fn read_events_never_touches_the_file() {
let (_tmp, path) = open_tmp();
{
let (mut journal, _) = Journal::open(&path).expect("open");
journal.append(|seq| user_message(seq, "kept"));
}
let mut raw = std::fs::read_to_string(&path).expect("read");
raw.push_str(r#"{"v":1,"seq":2,"at":"2026-"#);
std::fs::write(&path, &raw).expect("corrupt");
let events = read_events(&path);
assert_eq!(events.len(), 1);
assert_eq!(std::fs::read_to_string(&path).expect("reread"), raw);
let ghost = path.parent().unwrap().join("ghost.jsonl");
assert!(read_events(&ghost).is_empty());
assert!(!ghost.exists());
}
#[test]
fn unknown_format_version_is_an_error() {
let (_tmp, path) = open_tmp();
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(
&path,
"{\"v\":2,\"seq\":1,\"at\":\"2026-07-04T00:00:00Z\",\"kind\":{\"type\":\"memory_added\",\"fact\":\"x\"}}\n",
)
.unwrap();
assert!(Journal::open(&path).is_err());
}
#[test]
fn event_round_trips_every_kind() {
let kinds = vec![
user_message(1, "hi"),
EventKind::UserMessage {
id: MessageId("msg-9".into()),
op: MessageOp::Say,
text: "worker report: PR landed".into(),
from: Some("worker".into()),
},
EventKind::TurnStarted {
turn_id: "turn-2".into(),
answers: vec![MessageId("msg-1".into())],
body: None,
},
EventKind::TurnItem {
turn_id: "turn-2".into(),
item: ConversationItem::Message {
id: "text-0".into(),
text: "working on it".into(),
phase: None,
},
},
EventKind::TurnSteered {
turn_id: "turn-2".into(),
answers: vec![MessageId("msg-3".into())],
},
EventKind::TurnFinished {
turn_id: "turn-2".into(),
status: Lifecycle::Completed,
usage: Usage {
input_tokens: Some(10),
output_tokens: Some(5),
cache_read_tokens: None,
cost_usd: Some(0.01),
},
termination_reason: None,
},
EventKind::LoopState {
from: LoopState::Idle,
to: LoopState::Turning {
turn_id: "turn-2".into(),
},
reason: "turn opened".into(),
},
EventKind::RunObserved {
run_id: "run-1".into(),
session_id: "sess-1".into(),
flow: "design".into(),
task: "sketch the journal".into(),
},
EventKind::RunCompleted {
run_id: "run-1".into(),
outcome: WorkerOutcome::Completed,
summary: "landed".into(),
},
EventKind::TaskObserved {
observation: task_observation(),
},
EventKind::ChannelOpened {
name: "ship.148e0e02".into(),
run_id: "run-1".into(),
},
EventKind::MemoryUpdated {
summary: "learned the fold".into(),
},
EventKind::MemoryAdded {
fact: "learned a fact".into(),
},
EventKind::ServerStarted {
pid: 4242,
endpoint: "127.0.0.1:50123".into(),
},
];
for (i, kind) in kinds.into_iter().enumerate() {
let event = Event {
v: FORMAT_VERSION,
seq: i as u64 + 1,
at: OffsetDateTime::now_utc(),
kind,
};
let line = serde_json::to_string(&event).expect("serialize");
let decoded: Event = serde_json::from_str(&line).expect("deserialize");
assert_eq!(decoded, event);
}
}
#[test]
fn new_appends_carry_the_renamed_kind_names() {
let (_tmp, path) = open_tmp();
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
let (mut journal, _events) = Journal::open(&path).expect("journal opens");
journal.append(|_| EventKind::RunCompleted {
run_id: "run-2".into(),
outcome: WorkerOutcome::Failed,
summary: "process gone".into(),
});
let raw = std::fs::read_to_string(&path).unwrap();
let last = raw.lines().last().unwrap();
assert!(last.contains("\"run_completed\""), "{last}");
assert!(!last.contains("worker_finished"), "{last}");
}
#[test]
fn narration_renders_every_event_kind() {
let kinds = vec![
user_message(1, "hi"),
EventKind::TurnStarted {
turn_id: "turn-2".into(),
answers: vec![],
body: None,
},
EventKind::TurnItem {
turn_id: "turn-2".into(),
item: ConversationItem::Message {
id: "text-0".into(),
text: "working on it".into(),
phase: None,
},
},
EventKind::TurnItem {
turn_id: "turn-2".into(),
item: ConversationItem::Command {
id: "cmd-0".into(),
command: vec!["cargo".into(), "test".into()],
cwd: "/repo".into(),
status: Lifecycle::Completed,
output: None,
exit_code: Some(0),
duration_ms: Some(1200),
},
},
EventKind::TurnItem {
turn_id: "turn-2".into(),
item: ConversationItem::Thought {
id: "thought-0".into(),
text: "hmm".into(),
},
},
EventKind::TurnItem {
turn_id: "turn-2".into(),
item: ConversationItem::File {
id: "file-0".into(),
changes: vec![],
status: Lifecycle::Completed,
},
},
EventKind::TurnItem {
turn_id: "turn-2".into(),
item: ConversationItem::Tool {
id: "tool-0".into(),
name: "Bash".into(),
status: Lifecycle::Completed,
input: None,
output: None,
},
},
EventKind::TurnSteered {
turn_id: "turn-2".into(),
answers: vec![MessageId("msg-3".into())],
},
EventKind::TurnFinished {
turn_id: "turn-2".into(),
status: Lifecycle::Completed,
usage: Usage::empty(),
termination_reason: None,
},
EventKind::LoopState {
from: LoopState::Idle,
to: LoopState::Turning {
turn_id: "turn-2".into(),
},
reason: "turn opened".into(),
},
EventKind::RunObserved {
run_id: "run-1".into(),
session_id: "sess-1".into(),
flow: "design".into(),
task: "sketch the journal".into(),
},
EventKind::RunCompleted {
run_id: "run-1".into(),
outcome: WorkerOutcome::Completed,
summary: "landed".into(),
},
EventKind::TaskObserved {
observation: task_observation(),
},
EventKind::ChannelOpened {
name: "ship.148e0e02".into(),
run_id: "run-1".into(),
},
EventKind::ServerStarted {
pid: 4242,
endpoint: "127.0.0.1:50123".into(),
},
];
let mut narrator = Narrator::default();
for kind in &kinds {
let narration = narrator.render(kind);
assert!(!narration.line.is_empty(), "silent narration for {kind:?}");
}
}
#[test]
fn typed_task_observation_uses_the_existing_durable_consumption_fold() {
let observation = task_observation();
let observation_id = MessageId(observation.inbox_id());
let observed = Event {
v: FORMAT_VERSION,
seq: 1,
at: OffsetDateTime::now_utc(),
kind: EventKind::TaskObserved {
observation: observation.clone(),
},
};
let pending = fold_thread(std::slice::from_ref(&observed));
assert_eq!(pending.pending_messages.len(), 1);
assert_eq!(pending.tasks.get(&observation_id), Some(&observation));
assert_eq!(pending.turns.len(), 1);
assert!(pending.turns[0].text.is_empty());
assert_eq!(
pending.turns[0]
.activity
.as_ref()
.map(|activity| activity.id.as_str()),
Some(observation_id.0.as_str())
);
let consumed = fold_thread(&[
observed,
Event {
v: FORMAT_VERSION,
seq: 2,
at: OffsetDateTime::now_utc(),
kind: EventKind::TurnStarted {
turn_id: "turn-2".to_string(),
answers: vec![observation_id],
body: None,
},
},
]);
assert!(consumed.pending_messages.is_empty());
}
#[test]
fn narration_demo_reads_like_a_console() {
let mut narrator = Narrator::default();
let mut render = |kind: EventKind| {
let narration = narrator.render(&kind);
println!(
"{:5} {}",
format!("{:?}", narration.level).to_uppercase(),
narration.line
);
narration
};
let n = render(EventKind::UserMessage {
id: MessageId("msg-1".into()),
op: MessageOp::Message,
text: "how is the reactive server refactor going?".into(),
from: None,
});
assert_eq!(n.level, NarrationLevel::Info);
assert_eq!(
n.line,
"chat ← \"how is the reactive server refactor going?\" (msg-1)"
);
let n = render(EventKind::UserMessage {
id: MessageId("msg-2".into()),
op: MessageOp::Say,
text: "run-42 landed: PR #12 merged, one clippy fix on the side".into(),
from: Some("worker".into()),
});
assert_eq!(
n.line,
"chat ← [worker] \"run-42 landed: PR #12 merged, one clippy fix on the side\" (msg-2)"
);
let n = render(EventKind::UserMessage {
id: MessageId("msg-3".into()),
op: MessageOp::Steer,
text: "focus on the journal tests first".into(),
from: None,
});
assert_eq!(
n.line,
"chat ← (steer) \"focus on the journal tests first\" (msg-3)"
);
let n = render(EventKind::LoopState {
from: LoopState::Idle,
to: LoopState::Turning {
turn_id: "turn-4".into(),
},
reason: "turn opened".into(),
});
assert_eq!(n.line, "state idle → turning (turn opened)");
let n = render(EventKind::TurnStarted {
turn_id: "turn-4".into(),
answers: vec![MessageId("msg-1".into()), MessageId("msg-2".into())],
body: None,
});
assert_eq!(n.line, "turn turn-4 opened (answers: msg-1, msg-2)");
let n = render(EventKind::TurnItem {
turn_id: "turn-4".into(),
item: ConversationItem::Message {
id: "text-0".into(),
text: "Checking the worker reports,\nthen answering the chat.".into(),
phase: None,
},
});
assert_eq!(n.level, NarrationLevel::Info);
assert_eq!(
n.line,
"loop: \"Checking the worker reports, then answering the chat.\""
);
let n = render(EventKind::TurnItem {
turn_id: "turn-4".into(),
item: ConversationItem::Command {
id: "cmd-0".into(),
command: vec!["git".into(), "log".into(), "--oneline".into(), "-5".into()],
cwd: "/repo".into(),
status: Lifecycle::Completed,
output: None,
exit_code: Some(0),
duration_ms: Some(80),
},
});
assert_eq!(n.line, " $ git log --oneline -5 → completed");
let n = render(EventKind::TurnItem {
turn_id: "turn-4".into(),
item: ConversationItem::Message {
id: "text-1".into(),
text: "The build worker is still grinding; I'll start the doc pass.".into(),
phase: None,
},
});
assert_eq!(n.level, NarrationLevel::Debug);
let n = render(EventKind::TurnItem {
turn_id: "turn-4".into(),
item: ConversationItem::Thought {
id: "thought-0".into(),
text: "the queue is empty after this".into(),
},
});
assert_eq!(n.level, NarrationLevel::Debug);
let n = render(EventKind::TurnSteered {
turn_id: "turn-4".into(),
answers: vec![MessageId("msg-3".into())],
});
assert_eq!(n.line, "turn turn-4 steered (answers: msg-3)");
let n = render(EventKind::TurnFinished {
turn_id: "turn-4".into(),
status: Lifecycle::Completed,
usage: Usage {
input_tokens: Some(192_400),
output_tokens: Some(1_400),
cache_read_tokens: Some(182_000),
cost_usd: Some(0.42),
},
termination_reason: None,
});
assert_eq!(
n.line,
"turn turn-4 completed · 4 items · 192k in / 1.4k out (182k cached)"
);
let n = render(EventKind::RunObserved {
run_id: "run-8c1d2e3f4a".into(),
session_id: "sess-9".into(),
flow: "build".into(),
task: "wire the narration tap into the journal".into(),
});
assert_eq!(
n.line,
"observed worker run-8c1d flow=build started · wire the narration tap into the journal"
);
let n = render(EventKind::RunCompleted {
run_id: "run-8c1d2e3f4a".into(),
outcome: WorkerOutcome::Completed,
summary: "narration tap landed, suite green".into(),
});
assert_eq!(
n.line,
"observed run run-8c1d completed · narration tap landed, suite green"
);
let n = render(EventKind::MemoryAdded {
fact: "workers report through the memory stream".into(),
});
assert_eq!(
n.line,
"memory added: workers report through the memory stream"
);
let n = render(EventKind::MemoryUpdated {
summary: "journal is the console's source of truth".into(),
});
assert_eq!(
n.line,
"memory curated: journal is the console's source of truth"
);
let n = render(EventKind::ServerStarted {
pid: 4242,
endpoint: "127.0.0.1:50123".into(),
});
assert_eq!(n.line, "server started · pid 4242 · 127.0.0.1:50123");
let n = render(EventKind::ChannelOpened {
name: "ship.148e0e02".into(),
run_id: "run-8c1d2e3f4a".into(),
});
assert_eq!(n.line, "channel ship.148e0e02 opened · run run-8c1d");
}
#[test]
fn fold_materializes_legacy_channel_opened_as_a_worker_turn() {
let events = vec![Event {
v: FORMAT_VERSION,
seq: 1,
at: OffsetDateTime::now_utc(),
kind: EventKind::ChannelOpened {
name: "ship.148e0e02".into(),
run_id: "run-7".into(),
},
}];
let fold = fold_thread(&events);
assert_eq!(fold.turns.len(), 1);
assert_eq!(fold.turns[0].text, "work line ship.148e0e02 opened");
assert_eq!(fold.turns[0].from.as_deref(), Some("worker"));
assert_eq!(fold.turns[0].id, "turn-1");
assert!(
fold.pending_messages.is_empty(),
"a channel opening never queues for the loop"
);
}
#[test]
fn narration_truncates_and_resets_per_turn() {
let mut narrator = Narrator::default();
let long = "x".repeat(100);
let n = narrator.render(&EventKind::UserMessage {
id: MessageId("msg-1".into()),
op: MessageOp::Message,
text: long.clone(),
from: None,
});
assert_eq!(n.line, format!("chat ← \"{}…\" (msg-1)", "x".repeat(60)));
for turn in ["turn-2", "turn-5"] {
narrator.render(&EventKind::TurnStarted {
turn_id: turn.into(),
answers: vec![],
body: None,
});
let n = narrator.render(&EventKind::TurnItem {
turn_id: turn.into(),
item: ConversationItem::Message {
id: "text-0".into(),
text: "gist".into(),
phase: None,
},
});
assert_eq!(n.level, NarrationLevel::Info, "each turn gets one gist");
narrator.render(&EventKind::TurnFinished {
turn_id: turn.into(),
status: Lifecycle::Completed,
usage: Usage::empty(),
termination_reason: None,
});
}
}
#[test]
fn fold_joins_message_items_into_text_and_keeps_state() {
let (_tmp, path) = open_tmp();
let (mut journal, _) = Journal::open(&path).expect("open");
let mut events = Vec::new();
events.push(journal.append(|seq| user_message(seq, "how goes it?")));
events.push(journal.append(|seq| EventKind::TurnStarted {
turn_id: format!("turn-{seq}"),
answers: vec![],
body: None,
}));
let turn_id = "turn-2".to_string();
events.push(journal.append(|_| EventKind::LoopState {
from: LoopState::Idle,
to: LoopState::Turning {
turn_id: turn_id.clone(),
},
reason: "turn opened".into(),
}));
events.push(journal.append(|_| EventKind::TurnItem {
turn_id: turn_id.clone(),
item: ConversationItem::Message {
id: "text-0".into(),
text: "first".into(),
phase: None,
},
}));
events.push(journal.append(|_| EventKind::TurnItem {
turn_id: turn_id.clone(),
item: ConversationItem::Tool {
id: "item-0".into(),
name: "Bash".into(),
status: Lifecycle::Completed,
input: None,
output: None,
},
}));
events.push(journal.append(|_| EventKind::TurnItem {
turn_id: turn_id.clone(),
item: ConversationItem::Message {
id: "text-1".into(),
text: "second".into(),
phase: None,
},
}));
events.push(journal.append(|_| EventKind::TurnFinished {
turn_id: turn_id.clone(),
status: Lifecycle::Completed,
usage: Usage::empty(),
termination_reason: None,
}));
let fold = fold_thread(&events);
assert!(fold.open.is_empty());
assert_eq!(fold.turns.len(), 2);
assert_eq!(fold.turns[0].role, ChatRole::User);
assert_eq!(fold.turns[0].id, "turn-1");
let assistant = &fold.turns[1];
assert_eq!(assistant.id, "turn-2");
assert_eq!(assistant.text, "first\nsecond");
assert_eq!(assistant.items.len(), 1);
assert_eq!(assistant.status, Lifecycle::Completed);
assert_eq!(
fold.state,
LoopState::Turning {
turn_id: turn_id.clone()
}
);
}
#[test]
fn fold_treats_say_as_attributed_consumable_input() {
let say = |seq: u64| EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op: MessageOp::Say,
text: "landed the parser PR".to_string(),
from: Some("worker".into()),
};
let events = vec![Event {
v: FORMAT_VERSION,
seq: 1,
at: OffsetDateTime::now_utc(),
kind: say(1),
}];
let fold = fold_thread(&events);
assert_eq!(fold.turns.len(), 1);
assert_eq!(fold.turns[0].role, ChatRole::User);
assert_eq!(fold.turns[0].from.as_deref(), Some("worker"));
assert_eq!(fold.pending_messages.len(), 1, "say queues for the loop");
assert_eq!(fold.pending_messages[0].op, MessageOp::Say);
assert_eq!(fold.pending_messages[0].from.as_deref(), Some("worker"));
let consumed = vec![
events[0].clone(),
Event {
v: FORMAT_VERSION,
seq: 2,
at: OffsetDateTime::now_utc(),
kind: EventKind::TurnStarted {
turn_id: "turn-2".into(),
answers: vec![MessageId("msg-1".into())],
body: None,
},
},
];
let fold = fold_thread(&consumed);
assert!(fold.pending_messages.is_empty(), "answered say is consumed");
}
#[test]
fn fold_surfaces_unfinished_turns_as_open() {
let (_tmp, path) = open_tmp();
let (mut journal, _) = Journal::open(&path).expect("open");
let started = journal.append(|seq| EventKind::TurnStarted {
turn_id: format!("turn-{seq}"),
answers: vec![],
body: None,
});
let turn_id = match &started.kind {
EventKind::TurnStarted { turn_id, .. } => turn_id.clone(),
_ => unreachable!(),
};
let item = journal.append(|_| EventKind::TurnItem {
turn_id,
item: ConversationItem::Message {
id: "text-0".into(),
text: "half a thought".into(),
phase: None,
},
});
let fold = fold_thread(&[started, item]);
assert!(fold.turns.is_empty());
assert_eq!(fold.open.len(), 1);
assert_eq!(fold.open[0].text, "half a thought");
assert_eq!(fold.open[0].status, Lifecycle::Running);
}
}