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::ProjectObservation;
use crate::task::TaskObservation;
use crate::wave::chat::{ChatBacking, ConversationEpoch};
use crate::wave::playhead::{BodyProvenance, Playhead, PlayheadEvent};
use crate::wave::state::LoopState;
use crate::wave::PromotionWake;
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,
}
#[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 source: Option<DiscordMessageSource>,
}
impl PendingMessage {
pub fn destination(&self) -> MessageDestination {
self.source
.as_ref()
.map(|source| MessageDestination::Discord(source.binding.clone()))
.unwrap_or(MessageDestination::Local)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MessageDestination {
Local,
Discord(DiscordChatBinding),
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct DiscordChatBinding {
pub guild_id: String,
pub channel_id: String,
}
impl DiscordChatBinding {
pub fn channel_url(&self) -> String {
format!(
"https://discord.com/channels/{}/{}",
self.guild_id, self.channel_id
)
}
pub fn message_url(&self, message_id: &str) -> String {
format!("{}/{}", self.channel_url(), message_id)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct DiscordMessageSource {
pub binding: DiscordChatBinding,
pub message_id: String,
pub author_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConversationEpochImport {
pub epoch: ConversationEpoch,
pub turn_ids: Vec<String>,
}
impl DiscordMessageSource {
pub fn uri(&self) -> String {
format!(
"discord://{}/{}/{}",
self.binding.guild_id, self.binding.channel_id, self.message_id
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DiscordMessagePart {
pub part_id: String,
pub nonce: String,
pub content: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DiscordDelivery {
pub delivery_id: String,
pub turn_id: String,
pub binding: DiscordChatBinding,
pub sources: Vec<DiscordMessageSource>,
pub parts: Vec<DiscordMessagePart>,
pub confirmed: HashMap<String, String>,
}
impl DiscordDelivery {
pub fn reply_to(&self) -> Option<&DiscordMessageSource> {
self.sources.last()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DiscordAttachment {
pub binding: DiscordChatBinding,
pub bot_user_id: String,
pub cursor: 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 {
ConversationEpochsImported {
epochs: Vec<ConversationEpochImport>,
},
ConversationEpochStarted {
epoch_id: String,
number: u64,
backing: ChatBacking,
},
UserMessage {
id: MessageId,
op: MessageOp,
text: String,
},
DiscordChatAttached {
binding: DiscordChatBinding,
bot_user_id: String,
cursor: Option<String>,
},
DiscordUserMessage {
id: MessageId,
text: String,
source: DiscordMessageSource,
},
DiscordAuthoredMessage {
id: MessageId,
op: MessageOp,
text: String,
source: DiscordMessageSource,
},
DiscordChatCursorAdvanced {
binding: DiscordChatBinding,
message_id: String,
},
DiscordChatSendPlanned {
delivery_id: String,
turn_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
binding: Option<DiscordChatBinding>,
sources: Vec<DiscordMessageSource>,
parts: Vec<DiscordMessagePart>,
},
DiscordChatSendConfirmed {
delivery_id: String,
part_id: String,
provider_message_id: 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,
},
PromotionObserved {
parent_wave_id: crate::id::WaveId,
parent: 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")
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReadOnlyJournalState {
Available,
Missing,
Partial,
Unavailable,
}
#[derive(Debug)]
pub struct ReadOnlyJournal {
pub events: Vec<Event>,
pub state: ReadOnlyJournalState,
pub detail: Option<String>,
}
pub fn read_events_with_state(path: &Path) -> ReadOnlyJournal {
let raw = match std::fs::read_to_string(path) {
Ok(raw) => raw,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
return ReadOnlyJournal {
events: Vec::new(),
state: ReadOnlyJournalState::Missing,
detail: Some("No durable Wave Chat history exists yet.".to_string()),
};
}
Err(err) => {
return ReadOnlyJournal {
events: Vec::new(),
state: ReadOnlyJournalState::Unavailable,
detail: Some(format!("Could not read durable Wave Chat history: {err}")),
};
}
};
let mut events = Vec::new();
for (index, line) in raw.lines().enumerate() {
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"
);
return ReadOnlyJournal {
events,
state: ReadOnlyJournalState::Partial,
detail: Some(format!(
"Durable Wave Chat history stops at line {}: format v{} is incompatible with v{FORMAT_VERSION}.",
index + 1,
event.v
)),
};
}
Err(_) => {
return ReadOnlyJournal {
events,
state: ReadOnlyJournalState::Partial,
detail: Some(format!(
"Durable Wave Chat history stops at unreadable line {}.",
index + 1
)),
};
}
}
}
ReadOnlyJournal {
events,
state: ReadOnlyJournalState::Available,
detail: None,
}
}
pub fn read_events(path: &Path) -> Vec<Event> {
read_events_with_state(path).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::ConversationEpochsImported { epochs } => {
info(format!("{} legacy chat epoch(s) imported", epochs.len()))
}
EventKind::ConversationEpochStarted {
number, backing, ..
} => info(format!("chat epoch {number} started · {backing:?}")),
EventKind::UserMessage { id, op, text } => {
let op_tag = match op {
MessageOp::Message => "",
MessageOp::Steer => "(steer) ",
MessageOp::Interrupt => "(interrupt) ",
};
info(format!("chat ← {op_tag}\"{}\" ({id})", ellipsize(text, 60)))
}
EventKind::DiscordChatAttached {
binding, cursor, ..
} => info(format!(
"discord attached · channel {} · cursor {}",
binding.channel_id,
cursor.as_deref().unwrap_or("empty")
)),
EventKind::DiscordUserMessage { id, text, source } => info(format!(
"discord ← \"{}\" ({id}, {})",
ellipsize(text, 60),
source.message_id
)),
EventKind::DiscordAuthoredMessage {
id,
op,
text,
source,
} => info(format!(
"discord app ← ({op:?}) \"{}\" ({id}, {})",
ellipsize(text, 60),
source.message_id
)),
EventKind::DiscordChatCursorAdvanced { message_id, .. } => {
debug(format!("discord cursor → {message_id}"))
}
EventKind::DiscordChatSendPlanned {
delivery_id, parts, ..
} => info(format!(
"discord send planned · {delivery_id} · {} part(s)",
parts.len()
)),
EventKind::DiscordChatSendConfirmed {
delivery_id,
part_id,
provider_message_id,
} => info(format!(
"discord send confirmed · {delivery_id}/{part_id} → {provider_message_id}"
)),
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::PromotionObserved {
parent_wave_id,
parent,
} => info(format!(
"observed promotion from Wave {} · {}",
parent, parent_wave_id
)),
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(id: &str) -> String {
id.chars().take(8).collect()
}
#[derive(Debug)]
pub struct Journal {
file: File,
next_seq: u64,
narrator: Narrator,
#[cfg(test)]
next_append_failure: Option<JournalAppendStage>,
}
#[derive(Debug, thiserror::Error)]
#[error("journal append at seq {seq} failed during {operation}: {source}")]
pub struct JournalAppendError {
seq: u64,
operation: &'static str,
#[source]
source: std::io::Error,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum JournalAppendStage {
Write,
Flush,
}
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(),
#[cfg(test)]
next_append_failure: None,
},
events,
))
}
#[cfg(test)]
pub fn next_seq(&self) -> u64 {
self.next_seq
}
pub fn append(&mut self, build: impl FnOnce(u64) -> EventKind) -> Event {
self.try_append(build)
.expect("journal truth must persist before projecting an event in memory")
}
pub fn try_append(
&mut self,
build: impl FnOnce(u64) -> EventKind,
) -> Result<Event, JournalAppendError> {
let event = Event {
v: FORMAT_VERSION,
seq: self.next_seq,
at: OffsetDateTime::now_utc(),
kind: build(self.next_seq),
};
let mut line = serde_json::to_vec(&event).map_err(|source| JournalAppendError {
seq: event.seq,
operation: "serialization",
source: std::io::Error::new(std::io::ErrorKind::InvalidData, source),
})?;
line.push(b'\n');
let checkpoint = self
.file
.metadata()
.map_err(|source| JournalAppendError {
seq: event.seq,
operation: "file inspection",
source,
})?
.len();
let write_result = match self.take_injected_failure(JournalAppendStage::Write) {
Some(source) => Err(source),
None => self.file.write_all(&line),
};
if let Err(source) = write_result {
return Err(self.rollback(event.seq, checkpoint, "write", source));
}
let flush_result = match self.take_injected_failure(JournalAppendStage::Flush) {
Some(source) => Err(source),
None => self.file.flush(),
};
if let Err(source) = flush_result {
return Err(self.rollback(event.seq, checkpoint, "flush", source));
}
self.next_seq += 1;
self.narrator.narrate(&event.kind);
Ok(event)
}
fn rollback(
&mut self,
seq: u64,
checkpoint: u64,
operation: &'static str,
source: std::io::Error,
) -> JournalAppendError {
match self.file.set_len(checkpoint) {
Ok(()) => JournalAppendError {
seq,
operation,
source,
},
Err(rollback) => JournalAppendError {
seq,
operation: "write rollback",
source: std::io::Error::new(
source.kind(),
format!(
"{operation} failed: {source}; truncating to byte {checkpoint} also failed: {rollback}"
),
),
},
}
}
#[cfg(test)]
pub(crate) fn fail_next_append(&mut self, failure: JournalAppendStage) {
self.next_append_failure = Some(failure);
}
fn take_injected_failure(
&mut self,
#[cfg_attr(not(test), allow(unused_variables))] expected: JournalAppendStage,
) -> Option<std::io::Error> {
#[cfg(test)]
if self.next_append_failure == Some(expected) {
self.next_append_failure = None;
return Some(std::io::Error::other("injected journal append failure"));
}
None
}
}
#[derive(Debug)]
pub struct ThreadFold {
pub turns: Vec<ChatTurn>,
pub conversation_epochs: Vec<ConversationEpoch>,
pub conversation_epoch_turns: HashMap<String, Vec<String>>,
pub discord_turn_bindings: HashMap<String, DiscordChatBinding>,
pub open: Vec<ChatTurn>,
pub state: LoopState,
pub playhead: Option<Playhead>,
pub pending_messages: Vec<PendingMessage>,
pub messages: HashMap<MessageId, PendingMessage>,
pub discord: Option<DiscordAttachment>,
pub discord_deliveries: HashMap<String, DiscordDelivery>,
pub completed_claims: HashMap<String, Vec<MessageId>>,
pub tasks: HashMap<MessageId, TaskObservation>,
pub projects: HashMap<MessageId, ProjectObservation>,
pub(crate) promotions: HashMap<MessageId, PromotionWake>,
pub open_claims: Vec<MessageId>,
}
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
}
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 conversation_epochs: Vec<ConversationEpoch> = Vec::new();
let mut conversation_epoch_turns: HashMap<String, Vec<String>> = HashMap::new();
let mut discord_turn_bindings: HashMap<String, DiscordChatBinding> = HashMap::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 discord: Option<DiscordAttachment> = None;
let mut discord_deliveries: HashMap<String, DiscordDelivery> = HashMap::new();
let mut completed_claims: HashMap<String, Vec<MessageId>> = HashMap::new();
let mut tasks: HashMap<MessageId, TaskObservation> = HashMap::new();
let mut projects: HashMap<MessageId, ProjectObservation> = HashMap::new();
let mut promotions: HashMap<MessageId, PromotionWake> = HashMap::new();
let mut consumed_messages: HashSet<MessageId> = HashSet::new();
let mut claims_by_open_turn: HashMap<String, Vec<MessageId>> = HashMap::new();
for event in events {
match &event.kind {
EventKind::ConversationEpochsImported { epochs } => {
for imported in epochs {
if let Some(active) = conversation_epochs.last_mut() {
active.ended_at = Some(imported.epoch.started_at.clone());
}
conversation_epoch_turns
.insert(imported.epoch.id.clone(), imported.turn_ids.clone());
conversation_epochs.push(imported.epoch.clone());
}
}
EventKind::ConversationEpochStarted {
epoch_id,
number,
backing,
} => {
let at = event.at_rfc3339();
if let Some(active) = conversation_epochs.last_mut() {
active.ended_at = Some(at.clone());
}
conversation_epochs.push(ConversationEpoch {
id: epoch_id.clone(),
number: *number,
backing: backing.clone(),
journal_seq: event.seq,
started_at: at,
ended_at: None,
});
}
EventKind::UserMessage { id, op, text } => {
let mut turn = ChatTurn::user(format!("turn-{}", event.seq), text.clone());
turn.created_at = event.at_rfc3339();
turns.push(turn);
let message = PendingMessage {
id: id.clone(),
op: *op,
text: text.clone(),
source: None,
};
if !consumed_messages.contains(id) {
pending_messages.push(message.clone());
}
messages.insert(id.clone(), message);
}
EventKind::DiscordChatAttached {
binding,
bot_user_id,
cursor,
} => {
discord = Some(DiscordAttachment {
binding: binding.clone(),
bot_user_id: bot_user_id.clone(),
cursor: cursor.clone(),
});
}
EventKind::DiscordUserMessage { id, text, source } => {
let turn_id = format!("turn-{}", event.seq);
let mut turn = ChatTurn::user(turn_id.clone(), text.clone());
turn.created_at = event.at_rfc3339();
turns.push(turn);
discord_turn_bindings.insert(turn_id, source.binding.clone());
let message = PendingMessage {
id: id.clone(),
op: MessageOp::Message,
text: format!("[{}]\n{}", source.uri(), text),
source: Some(source.clone()),
};
if !consumed_messages.contains(id) {
pending_messages.push(message.clone());
}
messages.insert(id.clone(), message);
}
EventKind::DiscordAuthoredMessage {
id,
op,
text,
source,
} => {
let turn_id = format!("turn-{}", event.seq);
let mut turn = ChatTurn::user(turn_id.clone(), text.clone());
turn.created_at = event.at_rfc3339();
turns.push(turn);
discord_turn_bindings.insert(turn_id, source.binding.clone());
let message = PendingMessage {
id: id.clone(),
op: *op,
text: format!("[{}]\n{}", source.uri(), text),
source: Some(source.clone()),
};
if !consumed_messages.contains(id) {
pending_messages.push(message.clone());
}
messages.insert(id.clone(), message);
}
EventKind::DiscordChatCursorAdvanced {
binding,
message_id,
} => {
if let Some(attached) = discord.as_mut() {
if &attached.binding == binding {
attached.cursor = Some(message_id.clone());
}
}
}
EventKind::DiscordChatSendPlanned {
delivery_id,
turn_id,
binding,
sources,
parts,
} => {
let destination = binding
.clone()
.or_else(|| sources.first().map(|source| source.binding.clone()));
if let Some(binding) = destination {
discord_deliveries
.entry(delivery_id.clone())
.or_insert_with(|| DiscordDelivery {
delivery_id: delivery_id.clone(),
turn_id: turn_id.clone(),
binding,
sources: sources.clone(),
parts: parts.clone(),
confirmed: HashMap::new(),
});
} else {
tracing::warn!(%delivery_id, "Discord delivery has no destination; ignoring it");
}
}
EventKind::DiscordChatSendConfirmed {
delivery_id,
part_id,
provider_message_id,
} => {
if let Some(delivery) = discord_deliveries.get_mut(delivery_id) {
delivery
.confirmed
.insert(part_id.clone(), provider_message_id.clone());
}
}
EventKind::TaskObserved { observation } => {
let message = task_observation_message(observation);
let turn = ChatTurn::child_activity(
format!("turn-{}", event.seq),
event.at_rfc3339(),
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(),
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::PromotionObserved {
parent_wave_id,
parent,
} => {
let wake = PromotionWake {
parent_wave_id: parent_wave_id.clone(),
parent: parent.clone(),
};
let message = promotion_wake_message(&wake);
if !consumed_messages.contains(&message.id) {
pending_messages.push(message.clone());
}
promotions.insert(message.id.clone(), wake);
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(),
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());
let claims = claims_by_open_turn.remove(turn_id).unwrap_or_default();
if *status == Lifecycle::Completed {
completed_claims.insert(turn_id.clone(), claims);
}
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::RunCompleted {
run_id,
outcome,
summary,
} => {
turns.push(run_completed_turn(event, run_id, *outcome, summary));
}
EventKind::RunObserved { .. } | EventKind::ServerStarted { .. } => {}
}
}
let open_claims = open
.iter()
.filter_map(|turn| claims_by_open_turn.remove(&turn.id))
.flatten()
.collect();
for delivery in discord_deliveries.values() {
if !delivery.confirmed.is_empty() {
discord_turn_bindings.insert(delivery.turn_id.clone(), delivery.binding.clone());
}
}
ThreadFold {
turns,
conversation_epochs,
conversation_epoch_turns,
discord_turn_bindings,
open,
state,
playhead,
pending_messages,
messages,
discord,
discord_deliveries,
completed_claims,
tasks,
projects,
promotions,
open_claims,
}
}
pub fn task_observation_message(observation: &TaskObservation) -> PendingMessage {
PendingMessage {
id: MessageId(observation.inbox_id()),
op: MessageOp::Message,
text: observation.prompt(),
source: None,
}
}
pub fn project_observation_message(observation: &ProjectObservation) -> PendingMessage {
PendingMessage {
id: MessageId(observation.inbox_id()),
op: MessageOp::Message,
text: observation.prompt(),
source: None,
}
}
pub(crate) fn promotion_wake_message(wake: &PromotionWake) -> PendingMessage {
PendingMessage {
id: MessageId(wake.inbox_id()),
op: MessageOp::Message,
text: wake.prompt(),
source: None,
}
}
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(),
}
}
fn task_observation() -> crate::task::TaskObservation {
crate::task::TaskObservation {
task_id: crate::task::TaskId::from_raw("task_example"),
issue_identifier: "INF-123".to_string(),
event_id: 42,
event: crate::task::TaskEventKind::Failed {
error: "provider stopped".to_string(),
resumable: true,
},
}
}
fn discord_binding() -> DiscordChatBinding {
DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
}
}
fn local_epoch() -> ConversationEpoch {
ConversationEpoch {
id: "chat-epoch-legacy-1".into(),
number: 1,
backing: ChatBacking::Local,
journal_seq: 0,
started_at: "2026-07-04T00:00:00Z".into(),
ended_at: None,
}
}
fn discord_source() -> DiscordMessageSource {
DiscordMessageSource {
binding: discord_binding(),
message_id: "101".into(),
author_id: "human".into(),
}
}
fn discord_parts() -> Vec<DiscordMessagePart> {
vec![DiscordMessagePart {
part_id: "part-1".into(),
nonce: "lf-nonce-1".into(),
content: "answer".into(),
}]
}
#[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 read = read_events_with_state(&path);
assert_eq!(read.events.len(), 1);
assert_eq!(read.state, ReadOnlyJournalState::Partial);
assert!(read.detail.as_deref().unwrap().contains("line 2"));
assert_eq!(std::fs::read_to_string(&path).expect("reread"), raw);
let ghost = path.parent().unwrap().join("ghost.jsonl");
let missing = read_events_with_state(&ghost);
assert!(missing.events.is_empty());
assert_eq!(missing.state, ReadOnlyJournalState::Missing);
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\":\"server_started\",\"pid\":1,\"endpoint\":\"127.0.0.1:1\"}}\n",
)
.unwrap();
assert!(Journal::open(&path).is_err());
}
#[test]
fn event_round_trips_every_kind() {
let kinds = vec![
user_message(1, "hi"),
EventKind::ConversationEpochsImported {
epochs: vec![ConversationEpochImport {
epoch: local_epoch(),
turn_ids: vec!["turn-1".into()],
}],
},
EventKind::ConversationEpochStarted {
epoch_id: "chat-epoch-2".into(),
number: 2,
backing: ChatBacking::Local,
},
EventKind::DiscordChatAttached {
binding: discord_binding(),
bot_user_id: "bot".into(),
cursor: Some("100".into()),
},
EventKind::DiscordUserMessage {
id: MessageId("msg-2".into()),
text: "question".into(),
source: discord_source(),
},
EventKind::DiscordAuthoredMessage {
id: MessageId("msg-3".into()),
op: MessageOp::Steer,
text: "change course".into(),
source: discord_source(),
},
EventKind::DiscordChatCursorAdvanced {
binding: discord_binding(),
message_id: "101".into(),
},
EventKind::DiscordChatSendPlanned {
delivery_id: "delivery-1".into(),
turn_id: "turn-3".into(),
binding: None,
sources: vec![discord_source()],
parts: discord_parts(),
},
EventKind::DiscordChatSendConfirmed {
delivery_id: "delivery-1".into(),
part_id: "part-1".into(),
provider_message_id: "102".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::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 historical_sourced_delivery_derives_its_destination() {
let event = Event {
v: FORMAT_VERSION,
seq: 1,
at: OffsetDateTime::now_utc(),
kind: EventKind::DiscordChatSendPlanned {
delivery_id: "delivery-1".into(),
turn_id: "turn-1".into(),
binding: None,
sources: vec![discord_source()],
parts: discord_parts(),
},
};
let delivery = fold_thread(&[event])
.discord_deliveries
.remove("delivery-1")
.expect("historical delivery");
assert_eq!(delivery.binding, discord_binding());
}
#[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::ConversationEpochsImported {
epochs: vec![ConversationEpochImport {
epoch: local_epoch(),
turn_ids: vec!["turn-1".into()],
}],
},
EventKind::ConversationEpochStarted {
epoch_id: "chat-epoch-2".into(),
number: 2,
backing: ChatBacking::Local,
},
EventKind::DiscordChatAttached {
binding: discord_binding(),
bot_user_id: "bot".into(),
cursor: Some("100".into()),
},
EventKind::DiscordUserMessage {
id: MessageId("msg-2".into()),
text: "question".into(),
source: discord_source(),
},
EventKind::DiscordAuthoredMessage {
id: MessageId("msg-3".into()),
op: MessageOp::Steer,
text: "change course".into(),
source: discord_source(),
},
EventKind::DiscordChatCursorAdvanced {
binding: discord_binding(),
message_id: "101".into(),
},
EventKind::DiscordChatSendPlanned {
delivery_id: "delivery-1".into(),
turn_id: "turn-3".into(),
binding: Some(discord_binding()),
sources: vec![discord_source()],
parts: discord_parts(),
},
EventKind::DiscordChatSendConfirmed {
delivery_id: "delivery-1".into(),
part_id: "part-1".into(),
provider_message_id: "102".into(),
},
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::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(),
});
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-3".into()),
op: MessageOp::Steer,
text: "focus on the journal tests first".into(),
});
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::ServerStarted {
pid: 4242,
endpoint: "127.0.0.1:50123".into(),
});
assert_eq!(n.line, "server started · pid 4242 · 127.0.0.1:50123");
}
#[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(),
});
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_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);
}
}