use std::collections::HashMap;
use std::time::Duration;
use anyhow::Result;
use crate::chat::turns::{
ChatRole, ChatTurn, ChildActivityKind, ChildActivitySubject, ChildControlActivity,
};
use crate::chat::types::Lifecycle;
use crate::lf::commands::chat::{resolve_target, CliContext};
use crate::lf::WaveTargetArgs;
use crate::wave::journal::ellipsize;
use crate::wave::subscription::{stream_events, Frame};
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(),
context.env_channel.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;
renderer.render(frame);
})
.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,
finished: bool,
}
#[derive(Debug)]
struct Renderer {
turns: HashMap<String, TurnProgress>,
}
impl Renderer {
fn new() -> Self {
Self {
turns: HashMap::new(),
}
}
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(),
"memory" => vec![format!("memory curated: {}", ellipsize(&frame.data, 70))],
"memory-add" => vec![format!("memory added: {}", frame.data)],
"turn" => {
let Ok(turn) = serde_json::from_str::<ChatTurn>(&frame.data) else {
return vec![format!(
"(unparseable turn frame: {})",
ellipsize(&frame.data, 80)
)];
};
self.turn_lines(&turn)
}
_ => 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;
let who = turn.from.as_deref().unwrap_or("you");
lines.push(format!("{who} › {}", turn.text));
}
return lines;
}
progress.opened = true;
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 child_activity_line(activity: &ChildControlActivity) -> Option<String> {
match activity.kind {
ChildActivityKind::StateChanged
| ChildActivityKind::ControlApplied
| ChildActivityKind::Directed
| ChildActivityKind::Incorporated => return None,
ChildActivityKind::ControlUncertain
| ChildActivityKind::DecisionRequired
| ChildActivityKind::DecisionResolved
| ChildActivityKind::PullRequestOpened
| 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::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 turn_json(id: &str, role: &str, text: &str, status: &str, items: &str) -> String {
format!(
"{{\"id\":\"{id}\",\"role\":\"{role}\",\"text\":\"{text}\",\"status\":\"{status}\",\
\"items\":{items},\"created_at\":\"2026-07-04T00:00:00Z\",\"from\":null}}"
)
}
fn activity_turn_json(id: &str, kind: &str, title: &str, summary: &str) -> String {
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\",\
\"session_id\":\"ts_1\",\"kind\":\"{kind}\",\"title\":\"{title}\",\"summary\":\"{summary}\",\
\"directive_version\":null,\"command_id\":null,\"effect\":null,\"source\":null,\
\"decision_id\":null,\"options\":[]}}}}"
)
}
#[test]
fn sse_parser_splits_frames_on_blank_lines() {
let out = frames("event: state\ndata: idle\n\nevent: turn\ndata: {\"id\":1}\n\n");
assert_eq!(
out,
vec![
Frame {
event: "state".into(),
data: "idle".into()
},
Frame {
event: "turn".into(),
data: "{\"id\":1}".into()
},
]
);
}
#[test]
fn sse_parser_handles_crlf_comments_and_multiline_data() {
let out = frames(": ping\r\n\r\nevent: memory-add\r\ndata: first\r\ndata: second\r\n\r\n");
assert_eq!(
out,
vec![Frame {
event: "memory-add".into(),
data: "first\nsecond".into()
}]
);
assert!(frames("event: turn\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,
"turn",
turn_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,
"turn",
turn_json("turn-2", "assistant", "", "running", "[]"),
);
feed(
&mut renderer,
"turn",
turn_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,
"turn",
turn_json(
"turn-3",
"assistant",
"Projection landed. One test caught a stale signature; fixed and green.",
"completed",
&pursue_items,
),
);
assert_eq!(
out,
vec![
"you › make wave chat human-first",
"wave › The transcript renders every tool call as a card. I'll keep them out of the conversation.",
"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_and_memory_without_backend_state() {
let mut renderer = Renderer::new();
assert_eq!(
renderer.lines_for(&Frame {
event: "state".into(),
data: "turning".into()
}),
Vec::<String>::new()
);
assert_eq!(
renderer.lines_for(&Frame {
event: "memory-add".into(),
data: "workers report via lf radio pub with full detail".into()
}),
vec!["memory added: workers report via lf radio pub with full detail"]
);
let user = turn_json("turn-1", "user", "how goes it?", "completed", "[]");
assert_eq!(
renderer.lines_for(&Frame {
event: "turn".into(),
data: user.clone()
}),
vec!["you › how goes it?"]
);
assert!(renderer
.lines_for(&Frame {
event: "turn".into(),
data: user
})
.is_empty());
}
#[test]
fn conversation_shows_child_outcomes_but_not_lifecycle_churn() {
let mut renderer = Renderer::new();
let state = activity_turn_json(
"turn-1",
"state_changed",
"Task is running",
"provider turn is active",
);
assert!(renderer
.lines_for(&Frame {
event: "turn".into(),
data: state,
})
.is_empty());
let opened = activity_turn_json(
"turn-2",
"pull_request_opened",
"Opened PR #877",
"https://github.com/loopflowstudio/loopflow/pull/877",
);
assert_eq!(
renderer.lines_for(&Frame {
event: "turn".into(),
data: opened,
}),
vec!["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| turn_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: "turn".into(),
data: running("", "[]"),
});
assert!(lines.is_empty());
let lines = renderer.lines_for(&Frame {
event: "turn".into(),
data: running("thinking", "[]"),
});
assert_eq!(lines, vec!["wave › thinking"]);
let lines = renderer.lines_for(&Frame {
event: "turn".into(),
data: running("thinking", tool),
});
assert!(lines.is_empty());
let terminal = turn_json("turn-2", "assistant", "thinking", "completed", tool);
let lines = renderer.lines_for(&Frame {
event: "turn".into(),
data: terminal.clone(),
});
assert!(lines.is_empty());
assert!(renderer
.lines_for(&Frame {
event: "turn".into(),
data: terminal
})
.is_empty());
}
#[tokio::test]
async fn stream_events_delivers_replay_then_live_frames() {
use crate::wave::runtime::WaveRuntime;
use crate::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(
runtime.clone(),
server::ResidentDoor::new("test-token"),
None,
None,
server::ShutdownDoor::new(),
);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
for i in 0..20 {
runtime
.deliver(crate::wave::journal::MessageOp::Message, format!("old {i}"))
.expect("older user turn");
}
runtime
.deliver(crate::wave::journal::MessageOp::Message, "replayed".into())
.expect("user turn");
runtime
.append_memory("workers report via lf radio pub with full useful detail")
.unwrap();
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);
let _ = stream_events(&endpoint, "", &mut on_frame).await;
});
for _ in 0..200 {
if seen
.lock()
.unwrap()
.iter()
.filter(|f| f.event == "turn")
.count()
== 12
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
runtime
.append_memory("a fact published after subscribe")
.unwrap();
for _ in 0..200 {
if seen
.lock()
.unwrap()
.iter()
.any(|f| f.event == "memory-add" && f.data.contains("after subscribe"))
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
task.abort();
let frames = seen.lock().unwrap().clone();
assert_eq!(frames[0].event, "state", "replay opens with the state");
assert!(
frames
.iter()
.any(|f| f.event == "turn" && f.data.contains("replayed")),
"replayed turn arrives: {frames:?}"
);
assert!(
frames
.iter()
.all(|f| f.event != "turn" || !f.data.contains("old 0")),
"the default bounded replay omits older turns: {frames:?}"
);
assert_eq!(
frames.iter().filter(|f| f.event == "turn").count(),
12,
"human subscriptions replay the recent 12 turns"
);
assert!(
frames.iter().any(|f| {
f.event == "memory-add"
&& f.data == "workers report via lf radio pub with full useful detail"
}),
"replayed memory-add frame arrives with the full fact: {frames:?}"
);
assert!(
frames
.iter()
.any(|f| f.event == "memory-add" && f.data == "a fact published after subscribe"),
"live memory-add frame arrives: {frames:?}"
);
}
}