use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, OnceLock};
use std::time::Duration;
use serde::Deserialize;
use crate::chat::turns::{ChatRole, ChatTurn};
use crate::chat::types::{ConversationItem, Lifecycle};
use crate::wave::journal::{fold_thread, journal_path, read_events};
use crate::wave::server::endpoint_path;
pub const WAVE_CHAT_RECENT_TURNS: usize = 12;
pub const WAVE_CHAT_MAX_CHARS: usize = 4_000;
const LIVE_READ_TIMEOUT: Duration = Duration::from_secs(1);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AmbientWaveRef {
Id(String),
Name(String),
}
pub fn resolve_ambient_wave(
env_wave_id: Option<&str>,
repo_root: Option<&Path>,
) -> Option<AmbientWaveRef> {
if let Some(id) = env_wave_id.map(str::trim).filter(|value| !value.is_empty()) {
return Some(AmbientWaveRef::Id(id.to_string()));
}
let repo_root = repo_root?;
if !repo_git_info(repo_root).is_worktree_root {
return None;
}
crate::ops::util::resolve_wave_name(repo_root, None).map(AmbientWaveRef::Name)
}
pub fn resolve_ambient_channel(
env_channel: Option<&str>,
env_wave_id: Option<&str>,
repo_root: Option<&Path>,
) -> Option<AmbientWaveRef> {
if let Some(channel) = env_channel.map(str::trim).filter(|value| !value.is_empty()) {
return Some(AmbientWaveRef::Name(channel.to_string()));
}
resolve_ambient_wave(env_wave_id, repo_root)
}
pub fn resolve_ambient_channel_name(repo_root: &Path) -> Option<String> {
let env_channel = std::env::var(crate::lf::session::CHANNEL_ENV).ok();
let env_wave_id = std::env::var(crate::lf::session::WAVE_ID_ENV).ok();
match resolve_ambient_channel(
env_channel.as_deref(),
env_wave_id.as_deref(),
Some(repo_root),
)? {
AmbientWaveRef::Id(id) => wave_name_for_id(&id),
AmbientWaveRef::Name(name) => Some(name),
}
}
#[derive(Debug, Clone)]
struct RepoGitInfo {
is_worktree_root: bool,
origin: PathBuf,
}
fn repo_git_info(repo_root: &Path) -> RepoGitInfo {
static CACHE: OnceLock<Mutex<HashMap<PathBuf, RepoGitInfo>>> = OnceLock::new();
let cache = CACHE.get_or_init(Default::default);
let mut cache = cache
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
cache
.entry(repo_root.to_path_buf())
.or_insert_with(|| query_repo_git_info(repo_root))
.clone()
}
fn query_repo_git_info(repo_root: &Path) -> RepoGitInfo {
let not_a_root = || RepoGitInfo {
is_worktree_root: false,
origin: repo_root.to_path_buf(),
};
let Ok(output) = std::process::Command::new("git")
.arg("-C")
.arg(repo_root)
.args([
"rev-parse",
"--path-format=absolute",
"--show-toplevel",
"--git-common-dir",
])
.output()
else {
return not_a_root();
};
if !output.status.success() {
return not_a_root();
}
let stdout = String::from_utf8_lossy(&output.stdout);
let mut lines = stdout.lines().map(str::trim);
let toplevel = PathBuf::from(lines.next().unwrap_or_default());
let common_dir = PathBuf::from(lines.next().unwrap_or_default());
let toplevel = toplevel.canonicalize().unwrap_or(toplevel);
let root = repo_root
.canonicalize()
.unwrap_or_else(|_| repo_root.to_path_buf());
if toplevel != root {
return not_a_root();
}
RepoGitInfo {
is_worktree_root: true,
origin: common_dir
.parent()
.map(Path::to_path_buf)
.unwrap_or_else(|| repo_root.to_path_buf()),
}
}
fn wave_name_for_id(id: &str) -> Option<String> {
let id: crate::lfd::id::LfdId = id.parse().ok()?;
std::thread::spawn(move || {
let rt = tokio::runtime::Runtime::new().ok()?;
rt.block_on(async {
let store = crate::lfdb::open_existing_store().await?;
store.get_wave(&id).await.ok().flatten()
})
})
.join()
.ok()
.flatten()
.map(|wave| wave.name().to_string())
}
pub fn wave_origin(repo_root: &Path) -> PathBuf {
repo_git_info(repo_root).origin
}
pub fn gather_wave_chat(repo_root: &Path, wave: &str) -> Option<String> {
let origin = wave_origin(repo_root);
let turns = live_turns(&origin, wave).or_else(|| journal_turns(&origin, wave))?;
render_wave_chat(&turns)
}
const PARENT_CHAT_TURNS: usize = 6;
const PARENT_CHAT_MAX_CHARS: usize = 1_500;
const CHILD_CHAT_MAX_CHARS: usize = WAVE_CHAT_MAX_CHARS - PARENT_CHAT_MAX_CHARS;
pub fn gather_channel_chat(repo_root: &Path, channel: &str) -> Option<String> {
let wave = crate::wave::channel::family_head(channel);
if wave == channel {
return gather_wave_chat(repo_root, channel);
}
let own = {
let events = read_events(&journal_path(repo_root, channel));
let fold = fold_thread(&events);
let mut turns = fold.turns;
turns.extend(fold.open);
render_wave_chat_budget(&turns, WAVE_CHAT_RECENT_TURNS, CHILD_CHAT_MAX_CHARS)
};
let parent = {
let origin = wave_origin(repo_root);
live_turns(&origin, wave)
.or_else(|| journal_turns(&origin, wave))
.and_then(|turns| {
render_wave_chat_budget(&turns, PARENT_CHAT_TURNS, PARENT_CHAT_MAX_CHARS)
})
};
match (own, parent) {
(None, None) => None,
(own, parent) => {
let mut sections = Vec::new();
if let Some(own) = own {
sections.push(format!("## this work line ({channel})\n{own}"));
}
if let Some(parent) = parent {
sections.push(format!("## wave {wave}\n{parent}"));
}
Some(sections.join("\n\n"))
}
}
}
pub fn read_endpoint_pointer(origin: &Path, wave: &str) -> Option<String> {
let addr = std::fs::read_to_string(endpoint_path(origin, wave)).ok()?;
let addr = addr.trim();
(!addr.is_empty()).then(|| addr.to_string())
}
fn live_turns(origin: &Path, wave: &str) -> Option<Vec<ChatTurn>> {
#[derive(Debug, Deserialize)]
struct ConversationBody {
turns: Vec<ChatTurn>,
}
let addr = read_endpoint_pointer(origin, wave)?;
let body = http_get(format!(
"http://{addr}/conversation?limit={WAVE_CHAT_RECENT_TURNS}"
))?;
serde_json::from_str::<ConversationBody>(&body)
.ok()
.map(|payload| payload.turns)
}
fn journal_turns(origin: &Path, wave: &str) -> Option<Vec<ChatTurn>> {
let events = read_events(&journal_path(origin, wave));
if events.is_empty() {
return None;
}
let fold = fold_thread(&events);
let mut turns = fold.turns;
turns.extend(fold.open);
Some(turns)
}
fn http_get(url: String) -> Option<String> {
std::thread::spawn(move || {
let client = reqwest::blocking::Client::builder()
.connect_timeout(LIVE_READ_TIMEOUT)
.timeout(LIVE_READ_TIMEOUT)
.build()
.ok()?;
let response = client.get(url).send().ok()?;
if !response.status().is_success() {
return None;
}
response.text().ok()
})
.join()
.ok()
.flatten()
}
pub fn render_wave_chat(turns: &[ChatTurn]) -> Option<String> {
render_wave_chat_budget(turns, WAVE_CHAT_RECENT_TURNS, WAVE_CHAT_MAX_CHARS)
}
fn render_wave_chat_budget(
turns: &[ChatTurn],
max_turns: usize,
max_chars: usize,
) -> Option<String> {
let mut lines: Vec<String> = Vec::new();
let mut used = 0usize;
for turn in turns.iter().rev().take(max_turns) {
let mut line = render_turn_line(turn);
if line.is_empty() {
continue;
}
let cost = line.len() + usize::from(!lines.is_empty());
if used + cost > max_chars {
if lines.is_empty() {
truncate_on_char_boundary(&mut line, max_chars);
lines.push(line);
}
break;
}
used += cost;
lines.push(line);
}
if lines.is_empty() {
return None;
}
lines.reverse();
Some(lines.join("\n"))
}
fn render_turn_line(turn: &ChatTurn) -> String {
let speaker = turn.from.as_deref().unwrap_or(match turn.role {
ChatRole::User => "user",
ChatRole::Assistant => "wave",
});
let text = turn.text.trim();
let tool_items = turn
.items
.iter()
.filter(|item| !matches!(item, ConversationItem::Message { .. }))
.count();
if text.is_empty() && tool_items == 0 && turn.status == Lifecycle::Completed {
return String::new();
}
let mut line = format!("{speaker}: {text}");
if tool_items > 0 {
let plural = if tool_items == 1 { "" } else { "s" };
line.push_str(&format!(" ({tool_items} tool item{plural})"));
}
match turn.status {
Lifecycle::Failed => line.push_str(" [failed]"),
Lifecycle::Interrupted => line.push_str(" [interrupted]"),
Lifecycle::Running | Lifecycle::Pending => line.push_str(" [in progress]"),
Lifecycle::Completed => {}
}
line
}
fn truncate_on_char_boundary(value: &mut String, mut max: usize) {
if value.len() <= max {
return;
}
while max > 0 && !value.is_char_boundary(max) {
max -= 1;
}
value.truncate(max);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::wave::journal::{EventKind, Journal, MessageId, MessageOp, Usage};
use std::io::{Read, Write};
#[test]
fn ambient_rule_env_id_wins_and_is_trimmed() {
assert_eq!(
resolve_ambient_wave(Some(" wave-1 "), None),
Some(AmbientWaveRef::Id("wave-1".to_string()))
);
assert_eq!(resolve_ambient_wave(Some(" "), None), None);
assert_eq!(resolve_ambient_wave(None, None), None);
}
#[test]
fn ambient_rule_resolves_wave_worktree_but_not_nested_or_bare_dirs() {
let repo = loopflow_test_support::TestRepo::new();
let worktree = repo.create_wave_worktree("ship");
assert_eq!(
resolve_ambient_wave(None, Some(&worktree)),
Some(AmbientWaveRef::Name("ship".to_string()))
);
let nested = worktree.join("nested");
std::fs::create_dir_all(&nested).unwrap();
assert_eq!(resolve_ambient_wave(None, Some(&nested)), None);
let tmp = tempfile::tempdir().expect("tempdir");
assert_eq!(resolve_ambient_wave(None, Some(tmp.path())), None);
assert_eq!(
resolve_ambient_wave(Some("wave-1"), Some(&worktree)),
Some(AmbientWaveRef::Id("wave-1".to_string()))
);
}
#[test]
fn ambient_channel_env_wins_then_worktree_name() {
assert_eq!(
resolve_ambient_channel(Some(" ship.148e "), Some("wave-1"), None),
Some(AmbientWaveRef::Name("ship.148e".to_string()))
);
assert_eq!(
resolve_ambient_channel(Some(""), Some("wave-1"), None),
Some(AmbientWaveRef::Id("wave-1".to_string()))
);
let repo = loopflow_test_support::TestRepo::new();
let worktree = repo.create_wave_worktree("ship.148e");
assert_eq!(
resolve_ambient_channel(None, None, Some(&worktree)),
Some(AmbientWaveRef::Name("ship.148e".to_string())),
"the work-line worktree's name is its channel"
);
}
#[test]
fn channel_chat_in_a_work_line_carries_child_and_parent_sections() {
let repo = loopflow_test_support::TestRepo::new();
let worktree = repo.create_wave_worktree("goals.148e");
seed_journal(repo.path(), "goals", "wave-level question?");
let (mut child, _) =
Journal::open(&journal_path(&worktree, "goals.148e")).expect("child journal");
child.append(|seq| EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op: MessageOp::Say,
text: "child-level report".to_string(),
from: None,
});
drop(child);
let overlay =
gather_channel_chat(&worktree, "goals.148e").expect("work-line overlay renders");
assert!(overlay.contains("## this work line (goals.148e)"));
assert!(overlay.contains("child-level report"));
assert!(overlay.contains("## wave goals"));
assert!(overlay.contains("wave-level question?"));
let child_pos = overlay.find("child-level report").unwrap();
let parent_pos = overlay.find("wave-level question?").unwrap();
assert!(child_pos < parent_pos, "own channel first, parent after");
let wave_only = gather_channel_chat(repo.path(), "goals").expect("wave chat");
assert!(
!wave_only.contains("## "),
"no section headers at the wave home"
);
let bare = repo.create_wave_worktree("goals.bare0");
let overlay = gather_channel_chat(&bare, "goals.bare0").expect("parent-only overlay");
assert!(!overlay.contains("## this work line"));
assert!(overlay.contains("## wave goals"));
}
#[test]
fn wave_origin_of_a_worktree_is_the_main_checkout() {
let repo = loopflow_test_support::TestRepo::new();
let worktree = repo.create_wave_worktree("origin-check");
let origin = wave_origin(&worktree);
assert_eq!(
origin.canonicalize().unwrap(),
repo.path().canonicalize().unwrap()
);
let tmp = tempfile::tempdir().expect("tempdir");
assert_eq!(wave_origin(tmp.path()), tmp.path());
}
fn turn(role: ChatRole, text: &str) -> ChatTurn {
ChatTurn {
id: "turn-0".to_string(),
role,
text: text.to_string(),
status: Lifecycle::Completed,
items: Vec::new(),
created_at: "1970-01-01T00:00:00Z".to_string(),
from: None,
}
}
#[test]
fn render_notes_status_tools_and_attribution() {
let mut worker = turn(ChatRole::User, "worker report: PR landed");
worker.from = Some("worker-1".to_string());
let mut failed = turn(ChatRole::Assistant, "tried a build");
failed.status = Lifecycle::Failed;
failed.items.push(ConversationItem::Tool {
id: "t-0".to_string(),
name: "Bash".to_string(),
status: Lifecycle::Completed,
input: None,
output: None,
});
let rendered = render_wave_chat(&[worker, failed]).expect("chat renders");
assert_eq!(
rendered,
"worker-1: worker report: PR landed\nwave: tried a build (1 tool item) [failed]"
);
}
#[test]
fn render_budget_drops_oldest_first_and_keeps_newest() {
let mut turns = Vec::new();
for i in 0..30 {
turns.push(turn(
ChatRole::Assistant,
&format!("turn {i} {}", "x".repeat(600)),
));
}
let rendered = render_wave_chat(&turns).expect("chat renders");
assert!(rendered.len() <= WAVE_CHAT_MAX_CHARS);
assert!(rendered.contains("turn 29"), "newest turn survives");
assert!(!rendered.contains("turn 0 "), "oldest turns are dropped");
let first_kept = rendered.lines().next().unwrap();
let last_kept = rendered.lines().last().unwrap();
assert!(last_kept.contains("turn 29"));
assert!(first_kept < last_kept);
}
#[test]
fn render_caps_turn_count() {
let turns: Vec<ChatTurn> = (0..40)
.map(|i| turn(ChatRole::User, &format!("m{i}")))
.collect();
let rendered = render_wave_chat(&turns).expect("chat renders");
assert_eq!(rendered.lines().count(), WAVE_CHAT_RECENT_TURNS);
assert!(rendered.ends_with("user: m39"));
assert!(rendered.starts_with(&format!("user: m{}", 40 - WAVE_CHAT_RECENT_TURNS)));
}
#[test]
fn render_clips_a_single_oversized_newest_turn() {
let huge = turn(ChatRole::Assistant, &"y".repeat(WAVE_CHAT_MAX_CHARS * 2));
let rendered = render_wave_chat(&[huge]).expect("chat renders");
assert_eq!(rendered.len(), WAVE_CHAT_MAX_CHARS);
}
#[test]
fn render_empty_thread_is_none() {
assert!(render_wave_chat(&[]).is_none());
assert!(render_wave_chat(&[turn(ChatRole::Assistant, " ")]).is_none());
}
fn seed_journal(origin: &Path, wave: &str, text: &str) {
let (mut journal, _) = Journal::open(&journal_path(origin, wave)).expect("open journal");
journal.append(|seq| EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op: MessageOp::Message,
text: text.to_string(),
from: None,
});
journal.append(|seq| EventKind::TurnStarted {
turn_id: format!("turn-{seq}"),
answers: vec![MessageId("msg-1".to_string())],
});
journal.append(|_| EventKind::TurnItem {
turn_id: "turn-2".to_string(),
item: ConversationItem::Message {
id: "text-0".to_string(),
text: "on it".to_string(),
phase: None,
},
});
journal.append(|_| EventKind::TurnFinished {
turn_id: "turn-2".to_string(),
status: Lifecycle::Completed,
usage: Usage::empty(),
});
}
#[test]
fn journal_fold_feeds_chat_when_no_server_answers() {
let tmp = tempfile::tempdir().expect("tempdir");
seed_journal(tmp.path(), "goals", "how goes the build?");
let chat = gather_wave_chat(tmp.path(), "goals").expect("journal-backed chat");
assert!(chat.contains("user: how goes the build?"));
assert!(chat.contains("wave: on it"));
}
#[test]
fn no_wave_state_yields_no_chat() {
let tmp = tempfile::tempdir().expect("tempdir");
assert!(gather_wave_chat(tmp.path(), "ghost").is_none());
}
fn spawn_canned_conversation(turn_text: &str) -> std::net::SocketAddr {
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind");
let addr = listener.local_addr().expect("addr");
let body = serde_json::json!({
"turns": [{
"id": "turn-1",
"role": "assistant",
"text": turn_text,
"status": "completed",
"items": [],
"created_at": "1970-01-01T00:00:00Z",
"from": null,
}]
})
.to_string();
std::thread::spawn(move || {
if let Ok((mut socket, _)) = listener.accept() {
let mut request = Vec::new();
let mut buf = [0u8; 1024];
while !request.windows(4).any(|window| window == b"\r\n\r\n") {
match socket.read(&mut buf) {
Ok(0) | Err(_) => break,
Ok(read) => request.extend_from_slice(&buf[..read]),
}
}
let _ = write!(
socket,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
let _ = socket.shutdown(std::net::Shutdown::Write);
let _ = socket.read(&mut buf);
}
});
addr
}
#[test]
fn live_server_is_preferred_over_the_journal_fold() {
let tmp = tempfile::tempdir().expect("tempdir");
seed_journal(tmp.path(), "goals", "from the journal");
let addr = spawn_canned_conversation("from the live server");
crate::wave::server::write_endpoint(tmp.path(), "goals", addr).expect("endpoint");
let chat = gather_wave_chat(tmp.path(), "goals").expect("live chat");
assert!(chat.contains("from the live server"));
assert!(!chat.contains("from the journal"));
}
#[test]
fn dead_endpoint_falls_back_to_the_journal() {
let tmp = tempfile::tempdir().expect("tempdir");
seed_journal(tmp.path(), "goals", "from the journal");
let dead = std::net::TcpListener::bind("127.0.0.1:0").expect("bind");
let dead_addr = dead.local_addr().expect("addr");
drop(dead);
crate::wave::server::write_endpoint(tmp.path(), "goals", dead_addr).expect("endpoint");
let chat = gather_wave_chat(tmp.path(), "goals").expect("journal fallback");
assert!(chat.contains("from the journal"));
}
}