use std::collections::HashMap;
use std::ops::ControlFlow;
use std::time::Duration;
use anyhow::Result;
use crate::chat::turns::{
ChatRole, ChatTurn, ChildActivityKind, ChildActivitySubject, ChildControlActivity, TurnDelta,
};
use crate::chat::types::{ConversationItem, Lifecycle};
use crate::controller::wave::chat::{
ChatBacking, ChatMessageSource, ConversationEpoch, WaveChatMessage,
};
use crate::controller::wave::journal::ellipsize;
use crate::controller::wave::subscription::{stream_events, Frame};
use crate::lf::commands::chat::{resolve_target, CliContext};
use crate::lf::WaveTargetArgs;
const BACKOFF_FLOOR: Duration = Duration::from_secs(1);
const BACKOFF_CEIL: Duration = Duration::from_secs(30);
pub(crate) async fn follow(wave: Option<&str>) -> Result<()> {
let target = WaveTargetArgs {
wave: wave.map(str::to_string),
parent: false,
};
let mut renderer = Renderer::new();
let mut backoff = BACKOFF_FLOOR;
let mut waiting_note_shown = false;
loop {
let context = CliContext::detect().await;
let resolved = resolve_target(
&target,
context.store.as_ref(),
context.repo.as_deref(),
context.env_wave_id.as_deref(),
)
.await?;
let Some(resolved) = resolved else {
eprintln!("no wave here; nothing to follow");
return Ok(());
};
if let Some(endpoint) = &resolved.endpoint {
let mut saw_frame = false;
let result = stream_events(endpoint, "", &mut |frame| {
saw_frame = true;
let resync = frame.event == "resync";
renderer.render(frame);
if resync {
ControlFlow::Break(())
} else {
ControlFlow::Continue(())
}
})
.await;
if let Err(err) = result {
tracing::debug!(error = %err, wave = resolved.name, "event stream dropped");
}
if saw_frame {
backoff = BACKOFF_FLOOR;
}
waiting_note_shown = false;
} else if !waiting_note_shown {
eprintln!(
"wave '{}' has no live listener; waiting (start one with `lf wave {}`)",
resolved.name, resolved.name
);
waiting_note_shown = true;
}
tokio::time::sleep(backoff).await;
backoff = (backoff * 2).min(BACKOFF_CEIL);
}
}
#[derive(Debug, Default)]
struct TurnProgress {
opened: bool,
text_chars: usize,
items_shown: usize,
finished: bool,
}
#[derive(Debug)]
struct Renderer {
turns: HashMap<String, TurnProgress>,
open: HashMap<String, WaveChatMessage>,
epoch: Option<ConversationEpoch>,
}
impl Renderer {
fn new() -> Self {
Self {
turns: HashMap::new(),
open: HashMap::new(),
epoch: None,
}
}
fn render(&mut self, frame: Frame) {
for line in self.lines_for(&frame) {
println!("{line}");
}
}
fn lines_for(&mut self, frame: &Frame) -> Vec<String> {
match frame.event.as_str() {
"state" => Vec::new(),
"epoch" => {
let Ok(epoch) = serde_json::from_str::<ConversationEpoch>(&frame.data) else {
return vec![format!(
"(unparseable epoch frame: {})",
ellipsize(&frame.data, 80)
)];
};
if self.epoch.as_ref().map(|current| ¤t.id) == Some(&epoch.id) {
return Vec::new();
}
self.open.clear();
self.epoch = Some(epoch.clone());
vec![format!(
"chat · epoch {} · {}",
epoch.number,
backing_name(&epoch.backing)
)]
}
"message" => {
let Ok(message) = serde_json::from_str::<WaveChatMessage>(&frame.data) else {
return vec![format!(
"(unparseable message frame: {})",
ellipsize(&frame.data, 80)
)];
};
self.open.insert(message.turn.id.clone(), message.clone());
prefix_source(self.turn_lines(&message.turn), &message.source)
}
"message-delta" => {
let Ok(delta) = serde_json::from_str::<TurnDelta>(&frame.data) else {
return vec![format!(
"(unparseable message-delta frame: {})",
ellipsize(&frame.data, 80)
)];
};
let Some(message) = self.open.get_mut(&delta.turn_id) else {
return Vec::new();
};
message.turn.absorb_item(delta.item);
let turn = message.turn.clone();
let source = message.source.clone();
prefix_source(self.turn_lines(&turn), &source)
}
"resync" => {
self.open.clear();
Vec::new()
}
_ => Vec::new(),
}
}
fn turn_lines(&mut self, turn: &ChatTurn) -> Vec<String> {
let progress = self.turns.entry(turn.id.clone()).or_default();
if progress.finished {
return Vec::new();
}
let mut lines = Vec::new();
if let Some(activity) = &turn.activity {
if !progress.opened {
progress.opened = true;
progress.finished = true;
if let Some(line) = child_activity_line(activity) {
lines.push(line);
}
}
return lines;
}
if turn.role == ChatRole::User {
if !progress.opened {
progress.opened = true;
progress.finished = true;
lines.push(format!("you › {}", turn.text));
}
return lines;
}
progress.opened = true;
if turn.items.len() > progress.items_shown {
for item in &turn.items[progress.items_shown..] {
if let ConversationItem::Message { text, phase, .. } = item {
if phase.as_deref() == Some("commentary") {
for fragment in text.split('\n').filter(|f| !f.trim().is_empty()) {
lines.push(format!("wave › {fragment}"));
}
}
}
}
progress.items_shown = turn.items.len();
}
if turn.text.chars().count() > progress.text_chars {
let fresh: String = turn.text.chars().skip(progress.text_chars).collect();
progress.text_chars = turn.text.chars().count();
for fragment in fresh.split('\n').filter(|f| !f.trim().is_empty()) {
lines.push(format!("wave › {fragment}"));
}
}
if turn.status != Lifecycle::Running && turn.status != Lifecycle::Pending {
progress.finished = true;
if turn.status == Lifecycle::Failed {
lines.push("wave › Turn failed.".into());
} else if turn.status == Lifecycle::Interrupted {
lines.push("wave › Turn interrupted.".into());
}
}
lines
}
}
fn backing_name(backing: &ChatBacking) -> &'static str {
match backing {
ChatBacking::Local => "local",
ChatBacking::Discord { .. } => "discord",
}
}
fn prefix_source(lines: Vec<String>, source: &ChatMessageSource) -> Vec<String> {
let prefix = match source {
ChatMessageSource::Local { journal_seq } => format!("local event {journal_seq}"),
ChatMessageSource::Discord { .. } => "discord".to_string(),
};
lines
.into_iter()
.map(|line| format!("{prefix} · {line}"))
.collect()
}
fn child_activity_line(activity: &ChildControlActivity) -> Option<String> {
match activity.kind {
ChildActivityKind::StateChanged => return None,
ChildActivityKind::PrOpened | ChildActivityKind::Completed | ChildActivityKind::Failed => {}
}
let subject = match activity.subject {
ChildActivitySubject::Project => "project",
ChildActivitySubject::Task => "task",
};
let message = if activity.summary.is_empty() {
activity.title.clone()
} else {
format!("{} — {}", activity.title, activity.summary)
};
Some(format!("{subject} {} › {message}", activity.subject_id))
}
#[cfg(test)]
mod tests {
use super::*;
fn frames(raw: &str) -> Vec<Frame> {
use crate::controller::wave::subscription::SseFrameParser;
let mut parser = SseFrameParser::default();
let mut out = Vec::new();
for byte in raw.bytes() {
if let Some(frame) = parser.push(byte) {
out.push(frame);
}
}
out
}
fn message_json(id: &str, role: &str, text: &str, status: &str, items: &str) -> String {
let turn = format!(
"{{\"id\":\"{id}\",\"role\":\"{role}\",\"text\":\"{text}\",\"status\":\"{status}\",\
\"items\":{items},\"created_at\":\"2026-07-04T00:00:00Z\",\"from\":null}}"
);
format!(
"{{\"epoch_id\":\"chat-epoch-1\",\"source\":{{\"kind\":\"local\",\"journal_seq\":2}},\"turn\":{turn}}}"
)
}
fn activity_message_json(id: &str, kind: &str, title: &str, summary: &str) -> String {
let turn = format!(
"{{\"id\":\"{id}\",\"role\":\"user\",\"text\":\"\",\"status\":\"completed\",\
\"items\":[],\"created_at\":\"2026-07-04T00:00:00Z\",\"from\":\"task\",\
\"activity\":{{\"id\":\"activity-{id}\",\"subject\":\"task\",\"subject_id\":\"W2-132\",\
\"work_id\":\"ts_1\",\"kind\":\"{kind}\",\"title\":\"{title}\",\"summary\":\"{summary}\",\
\"directive_version\":null,\"command_id\":null,\"effect\":null,\"source\":null,\
\"decision_id\":null,\"options\":[]}}}}"
);
format!(
"{{\"epoch_id\":\"chat-epoch-1\",\"source\":{{\"kind\":\"local\",\"journal_seq\":2}},\"turn\":{turn}}}"
)
}
fn message_delta_json(turn_id: &str, item: &str) -> String {
format!("{{\"turn_id\":\"{turn_id}\",\"item\":{item}}}")
}
fn stream_message_item(id: &str, text: &str) -> String {
format!("{{\"type\":\"message\",\"id\":\"{id}\",\"text\":\"{text}\",\"phase\":\"stream\"}}")
}
fn commentary_message_item(id: &str, text: &str) -> String {
format!(
"{{\"type\":\"message\",\"id\":\"{id}\",\"text\":\"{text}\",\"phase\":\"commentary\"}}"
)
}
#[test]
fn source_bearing_messages_print_their_epoch_and_authority() {
let mut renderer = Renderer::new();
let epoch = ConversationEpoch {
id: "chat-epoch-1".into(),
number: 1,
backing: ChatBacking::Local,
journal_seq: 1,
started_at: "2026-07-04T00:00:00Z".into(),
ended_at: None,
};
assert_eq!(
renderer.lines_for(&Frame {
event: "epoch".into(),
data: serde_json::to_string(&epoch).expect("epoch JSON"),
}),
vec!["chat · epoch 1 · local"]
);
let local = WaveChatMessage {
epoch_id: epoch.id.clone(),
source: ChatMessageSource::Local { journal_seq: 7 },
turn: ChatTurn::user("turn-7".into(), "hello".into()),
};
assert_eq!(
renderer.lines_for(&Frame {
event: "message".into(),
data: serde_json::to_string(&local).expect("local message JSON"),
}),
vec!["local event 7 · you › hello"]
);
let mut assistant = ChatTurn::user("discord-9".into(), "shipped".into());
assistant.role = ChatRole::Assistant;
let discord = WaveChatMessage {
epoch_id: "chat-epoch-2".into(),
source: ChatMessageSource::Discord {
guild_id: "guild".into(),
channel_id: "channel".into(),
message_id: "9".into(),
author_id: "bot".into(),
url: "https://discord.com/channels/guild/channel/9".into(),
},
turn: assistant,
};
assert_eq!(
renderer.lines_for(&Frame {
event: "message".into(),
data: serde_json::to_string(&discord).expect("Discord message JSON"),
}),
vec!["discord · wave › shipped"]
);
}
#[test]
fn commentary_items_render_as_the_wave_speaking() {
let mut renderer = Renderer::new();
assert!(renderer
.lines_for(&Frame {
event: "message".into(),
data: message_json("turn-10", "assistant", "", "running", "[]"),
})
.is_empty());
assert_eq!(
renderer.lines_for(&Frame {
event: "message-delta".into(),
data: message_delta_json(
"turn-10",
&commentary_message_item("m-0", "Auditing the plan first."),
),
}),
vec!["local event 2 · wave › Auditing the plan first."]
);
assert_eq!(
renderer.lines_for(&Frame {
event: "message-delta".into(),
data: message_delta_json("turn-10", &stream_message_item("text-0", "Done.")),
}),
vec!["local event 2 · wave › Done."]
);
assert!(renderer
.lines_for(&Frame {
event: "message".into(),
data: message_json(
"turn-10",
"assistant",
"Done.",
"completed",
&format!(
"[{}]",
commentary_message_item("m-0", "Auditing the plan first.")
),
),
})
.is_empty());
}
#[test]
fn message_delta_frames_render_growth_and_resync_drops_reconstruction() {
let mut renderer = Renderer::new();
assert!(renderer
.lines_for(&Frame {
event: "message".into(),
data: message_json("turn-1", "assistant", "", "running", "[]"),
})
.is_empty());
assert_eq!(
renderer.lines_for(&Frame {
event: "message-delta".into(),
data: message_delta_json(
"turn-1",
&stream_message_item("text-0", "I fixed the parser.")
),
}),
vec!["local event 2 · wave › I fixed the parser."]
);
assert!(renderer
.lines_for(&Frame {
event: "message-delta".into(),
data: message_delta_json(
"turn-1",
"{\"type\":\"tool\",\"id\":\"t-1\",\"name\":\"Bash\",\"status\":\"completed\",\"input\":null,\"output\":null}",
),
})
.is_empty());
assert!(renderer
.lines_for(&Frame {
event: "resync".into(),
data: "reconnect".into(),
})
.is_empty());
assert!(
renderer
.lines_for(&Frame {
event: "message-delta".into(),
data: message_delta_json("turn-1", &stream_message_item("text-1", "more")),
})
.is_empty(),
"no reconstruction survives a resync; the whole-turn replay rebuilds it"
);
}
#[test]
fn sse_parser_splits_frames_on_blank_lines() {
let out = frames("event: state\ndata: idle\n\nevent: message\ndata: {\"id\":1}\n\n");
assert_eq!(
out,
vec![
Frame {
event: "state".into(),
data: "idle".into()
},
Frame {
event: "message".into(),
data: "{\"id\":1}".into()
},
]
);
}
#[test]
fn sse_parser_handles_crlf_comments_and_multiline_data() {
let out = frames(": ping\r\n\r\nevent: note\r\ndata: first\r\ndata: second\r\n\r\n");
assert_eq!(
out,
vec![Frame {
event: "note".into(),
data: "first\nsecond".into()
}]
);
assert!(frames("event: message\ndata: {\"id\":1}\n").is_empty());
}
#[test]
fn a_long_running_task_reads_as_conversation_not_a_build_log() {
let cmd = |id: &str, argv: &str, exit: i64| {
format!(
"{{\"type\":\"command\",\"id\":\"{id}\",\"command\":[\"sh\",\"-c\",\"{argv}\"],\
\"cwd\":\"/repo\",\"status\":\"completed\",\"output\":\"…\",\"exit_code\":{exit},\
\"duration_ms\":900}}"
)
};
let tool = |id: &str, name: &str| {
format!(
"{{\"type\":\"tool\",\"id\":\"{id}\",\"name\":\"{name}\",\"status\":\"completed\",\
\"input\":null,\"output\":\"ok\"}}"
)
};
let edit = |id: &str, path: &str| {
format!(
"{{\"type\":\"file\",\"id\":\"{id}\",\"changes\":[{{\"path\":\"{path}\",\
\"kind\":\"modified\",\"diff\":null}}],\"status\":\"completed\"}}"
)
};
let think = |id: &str| {
format!("{{\"type\":\"thought\",\"id\":\"{id}\",\"text\":\"weighing the options\"}}")
};
let mut renderer = Renderer::new();
let mut out = Vec::new();
let mut feed = |renderer: &mut Renderer, event: &str, data: String| {
out.extend(renderer.lines_for(&Frame {
event: event.into(),
data,
}));
};
feed(
&mut renderer,
"message",
message_json(
"turn-1",
"user",
"make wave chat human-first",
"completed",
"[]",
),
);
feed(&mut renderer, "state", "turning".into());
let clarify_items = format!(
"[{},{},{},{},{}]",
think("h0"),
tool("t0", "Read"),
tool("t1", "Grep"),
cmd("c0", "git log --oneline -3", 0),
edit("f0", "scratch/w2-129.md")
);
feed(
&mut renderer,
"message",
message_json("turn-2", "assistant", "", "running", "[]"),
);
feed(
&mut renderer,
"message",
message_json(
"turn-2",
"assistant",
"The transcript renders every tool call as a card. I'll keep them out of the conversation.",
"completed",
&clarify_items,
),
);
let pursue_items = format!(
"[{},{},{},{},{},{}]",
edit("f1", "swift/Loopflow/Models/WaveChatTranscript.swift"),
edit("f2", "swift/LoopflowMac/Views/MessageRow.swift"),
cmd("c1", "cargo build", 0),
cmd("c2", "swift test", 1),
tool("t2", "Edit"),
cmd("c3", "swift test", 0)
);
feed(
&mut renderer,
"message",
message_json(
"turn-3",
"assistant",
"Projection landed. One test caught a stale signature; fixed and green.",
"completed",
&pursue_items,
),
);
assert_eq!(
out,
vec![
"local event 2 · you › make wave chat human-first",
"local event 2 · wave › The transcript renders every tool call as a card. I'll keep them out of the conversation.",
"local event 2 · wave › Projection landed. One test caught a stale signature; fixed and green.",
],
"the conversation is prose; backend evidence remains in the journal"
);
}
#[test]
fn conversation_renders_chat_without_backend_state() {
let mut renderer = Renderer::new();
assert_eq!(
renderer.lines_for(&Frame {
event: "state".into(),
data: "turning".into()
}),
Vec::<String>::new()
);
let user = message_json("turn-1", "user", "how goes it?", "completed", "[]");
assert_eq!(
renderer.lines_for(&Frame {
event: "message".into(),
data: user.clone()
}),
vec!["local event 2 · you › how goes it?"]
);
assert!(renderer
.lines_for(&Frame {
event: "message".into(),
data: user
})
.is_empty());
}
#[test]
fn conversation_shows_child_outcomes_but_not_lifecycle_churn() {
let mut renderer = Renderer::new();
let state = activity_message_json(
"turn-1",
"state_changed",
"Task is running",
"provider turn is active",
);
assert!(renderer
.lines_for(&Frame {
event: "message".into(),
data: state,
})
.is_empty());
let opened = activity_message_json(
"turn-2",
"pr_opened",
"Opened PR #877",
"https://github.com/loopflowstudio/loopflow/pull/877",
);
assert_eq!(
renderer.lines_for(&Frame {
event: "message".into(),
data: opened,
}),
vec!["local event 2 · task W2-132 › Opened PR #877 — https://github.com/loopflowstudio/loopflow/pull/877"]
);
}
#[test]
fn conversation_prints_only_the_growth_of_a_repeating_turn_id() {
let mut renderer = Renderer::new();
let running =
|text: &str, items: &str| message_json("turn-2", "assistant", text, "running", items);
let tool = "[{\"type\":\"tool\",\"id\":\"t0\",\"name\":\"Bash\",\"status\":\"completed\",\
\"input\":null,\"output\":null}]";
let lines = renderer.lines_for(&Frame {
event: "message".into(),
data: running("", "[]"),
});
assert!(lines.is_empty());
let lines = renderer.lines_for(&Frame {
event: "message".into(),
data: running("thinking", "[]"),
});
assert_eq!(lines, vec!["local event 2 · wave › thinking"]);
let lines = renderer.lines_for(&Frame {
event: "message".into(),
data: running("thinking", tool),
});
assert!(lines.is_empty());
let terminal = message_json("turn-2", "assistant", "thinking", "completed", tool);
let lines = renderer.lines_for(&Frame {
event: "message".into(),
data: terminal.clone(),
});
assert!(lines.is_empty());
assert!(renderer
.lines_for(&Frame {
event: "message".into(),
data: terminal
})
.is_empty());
}
#[tokio::test]
async fn stream_events_delivers_replay_then_live_frames() {
use crate::controller::wave::runtime::WaveRuntime;
use crate::controller::wave::server;
use std::sync::{Arc, Mutex};
let tmp = tempfile::tempdir().expect("tempdir");
let runtime =
WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open runtime");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let app = server::router_with_observer(
runtime.clone(),
server::ResidentDoor::new("test-token"),
Arc::new(crate::controller::wave::registry::ObserverSlot::new(
runtime.clone(),
None,
)),
None,
server::ShutdownDoor::new(),
);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
for i in 0..20 {
runtime
.deliver(
crate::controller::wave::journal::MessageOp::Message,
format!("old {i}"),
)
.expect("older user turn");
}
runtime
.deliver(
crate::controller::wave::journal::MessageOp::Message,
"replayed".into(),
)
.expect("user turn");
let seen: Arc<Mutex<Vec<Frame>>> = Arc::new(Mutex::new(Vec::new()));
let sink = seen.clone();
let endpoint = addr.to_string();
let task = tokio::spawn(async move {
let mut on_frame = |frame: Frame| {
sink.lock().unwrap().push(frame);
ControlFlow::Continue(())
};
let _ = stream_events(&endpoint, "", &mut on_frame).await;
});
for _ in 0..200 {
if seen
.lock()
.unwrap()
.iter()
.filter(|f| f.event == "message")
.count()
== 12
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
runtime
.deliver(
crate::controller::wave::journal::MessageOp::Message,
"live after subscribe".into(),
)
.expect("live user turn");
for _ in 0..200 {
if seen
.lock()
.unwrap()
.iter()
.any(|f| f.event == "message" && f.data.contains("live after subscribe"))
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
task.abort();
let frames = seen.lock().unwrap().clone();
assert_eq!(frames[0].event, "epoch", "replay names its authority first");
assert_eq!(frames[1].event, "backing-health");
assert_eq!(frames[2].event, "state");
assert!(
frames
.iter()
.any(|f| f.event == "message" && f.data.contains("replayed")),
"replayed source-bearing message arrives: {frames:?}"
);
assert!(
frames
.iter()
.all(|f| f.event != "message" || !f.data.contains("old 0")),
"the default bounded replay omits older messages: {frames:?}"
);
assert_eq!(
frames.iter().filter(|f| f.event == "message").count(),
13,
"human subscriptions replay 12 messages, then stream the live message"
);
assert!(
frames
.iter()
.any(|f| f.event == "message" && f.data.contains("live after subscribe")),
"live source-bearing message arrives: {frames:?}"
);
}
}