use std::io::BufRead;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use anyhow::{anyhow, bail, Result};
use serde::{Deserialize, Serialize};
use crate::chat::turns::ChatTurn;
use crate::engine::wave_context::{
read_endpoint_pointer, resolve_ambient_channel, resolve_managed_wave_name, wave_origin,
AmbientChannelRef, WaveResolveError,
};
use crate::lf::commands::thread;
use crate::lf::commands::util::{find_repo_root, message_text};
use crate::lf::WaveTargetArgs;
use crate::store::{open_existing_store, SharedStore};
use crate::wave::channel::family_head;
use crate::wave::journal::{
fold_thread, journal_path, read_events_with_state, MessageOp, ReadOnlyJournalState,
};
use crate::wave::server::HUMAN_THREAD_REPLAY_LIMIT;
use crate::wave::Wave;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ChatHistoryState {
Available,
Missing,
Partial,
Unavailable,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ChatHistorySnapshot {
pub state: ChatHistoryState,
pub detail: Option<String>,
pub turns: Vec<ChatTurn>,
pub truncated: bool,
}
pub fn run(
text_args: &[String],
follow: bool,
steer: bool,
history: bool,
json: bool,
limit: Option<usize>,
target: &WaveTargetArgs,
) -> Result<()> {
let rt = tokio::runtime::Runtime::new()?;
rt.block_on(async {
let context = CliContext::detect().await;
if history {
if !json {
bail!("--history currently requires --json");
}
history_with_context(&context, target, limit.unwrap_or(HUMAN_THREAD_REPLAY_LIMIT)).await
} else if follow {
follow_with_context(&context, steer, target).await
} else {
run_with_context(&context, text_args, steer, target).await
}
})
}
async fn history_with_context(
context: &CliContext,
target: &WaveTargetArgs,
limit: usize,
) -> Result<()> {
if limit == 0 {
bail!("--limit must be at least 1");
}
let Some(resolved) = resolve_target(
target,
context.store.as_ref(),
context.repo.as_deref(),
context.env_wave_id.as_deref(),
context.env_channel.as_deref(),
)
.await?
else {
bail!("no wave here — name one with `lf chat --history --json -w <wave>`");
};
let repo_root = resolved.repo_root.ok_or_else(|| {
anyhow!(
"wave '{}' has no local origin repository for durable history",
resolved.name
)
})?;
let snapshot = history_snapshot(&repo_root, &resolved.name, limit);
println!("{}", serde_json::to_string(&snapshot)?);
Ok(())
}
fn history_snapshot(repo_root: &Path, wave: &str, limit: usize) -> ChatHistorySnapshot {
let read = read_events_with_state(&journal_path(repo_root, wave));
let state = match read.state {
ReadOnlyJournalState::Available => ChatHistoryState::Available,
ReadOnlyJournalState::Missing => ChatHistoryState::Missing,
ReadOnlyJournalState::Partial => ChatHistoryState::Partial,
ReadOnlyJournalState::Unavailable => ChatHistoryState::Unavailable,
};
let mut fold = fold_thread(&read.events);
fold.turns.append(&mut fold.open);
let truncated = fold.turns.len() > limit;
let keep_from = fold.turns.len().saturating_sub(limit);
let turns = fold.turns.split_off(keep_from);
ChatHistorySnapshot {
state,
detail: read.detail,
turns,
truncated,
}
}
pub(crate) async fn run_with_context(
context: &CliContext,
text_args: &[String],
steer: bool,
target: &WaveTargetArgs,
) -> Result<()> {
let Some(resolved) = resolve_target(
target,
context.store.as_ref(),
context.repo.as_deref(),
context.env_wave_id.as_deref(),
context.env_channel.as_deref(),
)
.await?
else {
eprintln!("no wave here; message dropped");
return Ok(());
};
let text = message_text(text_args, std::io::stdin())?;
let endpoint = resolved.require_endpoint()?;
post_message(&endpoint, &text, steer).await?;
println!(
"sent to '{}' ({})",
resolved.name,
if steer {
"steer live, otherwise queue"
} else {
"queued for the next turn"
}
);
Ok(())
}
async fn follow_with_context(
context: &CliContext,
steer: bool,
target: &WaveTargetArgs,
) -> Result<()> {
let Some(resolved) = resolve_target(
target,
context.store.as_ref(),
context.repo.as_deref(),
context.env_wave_id.as_deref(),
context.env_channel.as_deref(),
)
.await?
else {
bail!("no wave here — name one with `lf chat --follow -w <wave>`");
};
follow_thread(&resolved, steer).await
}
async fn follow_thread(resolved: &ResolvedWave, steer: bool) -> Result<()> {
let endpoint = resolved.require_endpoint()?;
println!(
"chat: {} @ {endpoint} (/help, Ctrl-D to leave)",
resolved.name
);
let wave_name = resolved.name.clone();
let stream = tokio::spawn(async move { thread::follow(Some(wave_name.as_str())).await });
let (tx, mut rx) = tokio::sync::mpsc::channel::<String>(16);
std::thread::spawn(move || {
for line in std::io::stdin().lock().lines() {
let Ok(line) = line else { break };
if tx.blocking_send(line).is_err() {
break;
}
}
});
loop {
let line = tokio::select! {
line = rx.recv() => line,
_ = tokio::signal::ctrl_c() => None,
};
let Some(line) = line else { break };
let line = line.trim();
if line.is_empty() {
continue;
}
if let Some(command) = line.strip_prefix('/') {
if !handle_command(command, &endpoint).await? {
break;
}
continue;
}
post_message(&endpoint, line, steer).await?;
}
stream.abort();
Ok(())
}
#[derive(Debug, PartialEq, Eq)]
enum Command {
Quit,
Status,
Help,
Unknown,
}
fn parse_command(command: &str) -> Command {
match command.trim() {
"q" | "quit" | "exit" => Command::Quit,
"status" => Command::Status,
"help" | "?" => Command::Help,
_ => Command::Unknown,
}
}
async fn handle_command(command: &str, endpoint: &str) -> Result<bool> {
match parse_command(command) {
Command::Quit => return Ok(false),
Command::Status => {
let health = get_json(endpoint, "/health").await?;
println!("{}", serde_json::to_string_pretty(&health)?);
}
Command::Help => println!(
" /status the wave's loop state\n \
/quit leave (Ctrl-D also works)\n \
anything else is spoken into the thread"
),
Command::Unknown => eprintln!("unknown command '/{}' — try /help", command.trim()),
}
Ok(true)
}
async fn post_message(endpoint: &str, text: &str, steer: bool) -> Result<()> {
let op = if steer {
MessageOp::Steer
} else {
MessageOp::Message
};
let body = serde_json::json!({ "op": op, "text": text });
post_json(endpoint, "/messages", &body).await?;
Ok(())
}
pub(crate) async fn nudge_child_observations(wave: &str) -> Result<bool> {
let context = CliContext::detect().await;
let target = WaveTargetArgs {
wave: Some(wave.to_string()),
parent: false,
};
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?
.ok_or_else(|| anyhow!("wave {wave:?} cannot be resolved"))?;
let Some(endpoint) = resolved.endpoint else {
return Ok(false);
};
post_json(&endpoint, "/observations", &serde_json::json!({})).await?;
Ok(true)
}
pub(crate) struct CliContext {
pub store: Option<SharedStore>,
pub repo: Option<PathBuf>,
pub env_wave_id: Option<String>,
pub env_channel: Option<String>,
}
impl CliContext {
pub async fn detect() -> Self {
Self {
store: open_existing_store().await.map(Arc::new),
repo: find_repo_root().ok(),
env_wave_id: std::env::var(crate::engine::wave_context::WAVE_ID_ENV)
.ok()
.filter(|value| !value.is_empty()),
env_channel: std::env::var(crate::engine::wave_context::CHANNEL_ENV)
.ok()
.filter(|value| !value.is_empty()),
}
}
}
#[derive(Debug)]
pub(crate) struct ResolvedWave {
pub name: String,
pub endpoint: Option<String>,
pub repo_root: Option<PathBuf>,
}
impl ResolvedWave {
pub fn require_endpoint(&self) -> Result<String> {
self.endpoint.clone().ok_or_else(|| {
anyhow!(
"wave '{name}' has no live listener — start one with `lf wave {name}`. \
(Queuing for offline waves is not implemented yet.)",
name = self.name
)
})
}
}
pub(crate) async fn resolve_target(
args: &WaveTargetArgs,
store: Option<&SharedStore>,
repo: Option<&Path>,
env_wave_id: Option<&str>,
env_channel: Option<&str>,
) -> Result<Option<ResolvedWave>> {
let main_repo = repo.map(wave_origin);
let mut own_row: Option<Wave> = None;
let mut own_name: Option<String> = None;
if args.wave.is_none() {
match resolve_ambient_channel(env_channel, env_wave_id) {
Some(AmbientChannelRef::WaveId(id)) => {
let name = resolve_managed_wave_name(store.map(|store| &**store), None, Some(&id))
.await
.map_err(|err| anyhow!("{err}"))?;
own_row = match store {
Some(store) => store.get_wave_by_name(&name).await?,
None => None,
};
own_name = Some(name);
}
Some(AmbientChannelRef::Channel(name)) => {
let head = family_head(&name).to_string();
if let Some(store) = store {
own_row = store.get_wave_by_name(&head).await?;
}
own_name = Some(head);
}
None => {}
}
}
let (target_row, target_name): (Option<Wave>, String) = if let Some(name) = &args.wave {
let head = family_head(name).to_string();
let store = store.ok_or_else(|| {
anyhow!(
"{}",
WaveResolveError::Registry("no wave registry on this machine".to_string())
)
})?;
let row = store
.get_wave_by_name(&head)
.await
.map_err(|err| anyhow!("failed to read wave registry: {err}"))?
.ok_or_else(|| anyhow!("{}", WaveResolveError::UnknownExplicit(head.clone())))?;
(Some(row), head)
} else if args.parent {
let store = store.ok_or_else(|| {
anyhow!(
"--parent needs the run registry to walk the wave tree, \
and this machine has none (~/.lf/loopflow.db)"
)
})?;
let own = own_row.ok_or_else(|| {
anyhow!("cannot resolve the invoking wave for --parent: no LF_WAVE_ID in env")
})?;
let parent = parent_wave(store, &own).await?;
let name = parent.name().to_string();
(Some(parent), name)
} else {
match own_row {
Some(row) => {
let name = row.name().to_string();
(Some(row), name)
}
None => {
let Some(name) = own_name.clone() else {
return Ok(None);
};
(None, name)
}
}
};
let mut endpoint = None;
let repo_root = target_row
.as_ref()
.map(|row| PathBuf::from(row.repo()))
.filter(|path| path.is_dir())
.or(main_repo);
if endpoint.is_none() {
if let Some(root) = &repo_root {
endpoint = read_endpoint_pointer(root, &target_name);
}
}
Ok(Some(ResolvedWave {
name: target_name,
endpoint,
repo_root,
}))
}
pub(crate) async fn parent_wave(store: &SharedStore, own: &Wave) -> Result<Wave> {
let parent_id = own.parent_wave_id().ok_or_else(|| {
anyhow!(
"wave '{}' has no parent — it is a root wave; the human \
fall-through arrives with Decisions",
own.name()
)
})?;
store.get_wave(parent_id).await?.ok_or_else(|| {
anyhow!(
"wave '{}' names parent {parent_id}, but the registry has no such wave",
own.name()
)
})
}
pub(crate) async fn post_json(
endpoint: &str,
path: &str,
body: &serde_json::Value,
) -> Result<serde_json::Value> {
let client = reqwest::Client::new();
let response = client
.post(format!("http://{endpoint}{path}"))
.json(body)
.send()
.await
.map_err(|err| {
anyhow!(
"wave listener at {endpoint} is not answering ({err}) — is `lf wave` still running?"
)
})?;
let status = response.status();
let text = response.text().await.unwrap_or_default();
if !status.is_success() {
bail!("wave server rejected the request ({status}): {text}");
}
Ok(serde_json::from_str(&text).unwrap_or(serde_json::Value::Null))
}
pub(crate) async fn get_json(endpoint: &str, path: &str) -> Result<serde_json::Value> {
let response = reqwest::get(format!("http://{endpoint}{path}"))
.await
.map_err(|err| {
anyhow!(
"wave listener at {endpoint} is not answering ({err}) — is `lf wave` still running?"
)
})?;
let status = response.status();
let text = response.text().await.unwrap_or_default();
if !status.is_success() {
bail!("wave server rejected the request ({status}): {text}");
}
serde_json::from_str(&text).map_err(|err| anyhow!("bad response from wave server: {err}"))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::lf::commands::fixtures::{boot_server, make_wave, temp_store};
use crate::wave::journal::{EventKind, MessageOp};
use crate::wave::runtime::InboxItem;
#[test]
fn durable_history_is_bounded_and_does_not_need_a_listener() {
let tmp = tempfile::tempdir().expect("tempdir");
let runtime =
crate::wave::runtime::WaveRuntime::open("ship".to_string(), tmp.path().to_path_buf())
.expect("open wave journal");
for index in 0..15 {
runtime
.deliver(MessageOp::Message, format!("message {index}"))
.expect("append message");
}
let snapshot = history_snapshot(tmp.path(), "ship", 12);
assert_eq!(snapshot.state, ChatHistoryState::Available);
assert_eq!(snapshot.turns.len(), 12);
assert!(snapshot.truncated);
assert_eq!(snapshot.turns.first().unwrap().text, "message 3");
assert_eq!(snapshot.turns.last().unwrap().text, "message 14");
}
#[test]
fn durable_history_distinguishes_missing_evidence_from_an_empty_thread() {
let tmp = tempfile::tempdir().expect("tempdir");
let snapshot = history_snapshot(tmp.path(), "ship", 12);
assert_eq!(snapshot.state, ChatHistoryState::Missing);
assert!(snapshot.turns.is_empty());
assert!(!snapshot.truncated);
assert!(snapshot.detail.is_some());
}
#[test]
fn chat_history_fixture_round_trips() {
let fixture = include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/../../tests/fixtures/dto/chat_history.json"
));
let snapshot: ChatHistorySnapshot =
serde_json::from_str(fixture).expect("decode chat history fixture");
assert_eq!(snapshot.state, ChatHistoryState::Partial);
assert_eq!(snapshot.turns.len(), 1);
assert!(snapshot.truncated);
let encoded = serde_json::to_string(&snapshot).expect("encode chat history fixture");
let decoded: ChatHistorySnapshot =
serde_json::from_str(&encoded).expect("re-decode chat history fixture");
assert_eq!(decoded, snapshot);
}
#[test]
fn interactive_commands_parse_and_everything_else_is_speech() {
assert_eq!(parse_command("quit"), Command::Quit);
assert_eq!(parse_command(" q "), Command::Quit);
assert_eq!(parse_command("exit"), Command::Quit);
assert_eq!(parse_command("status"), Command::Status);
assert_eq!(parse_command("help"), Command::Help);
assert_eq!(parse_command("?"), Command::Help);
assert_eq!(parse_command("deploy"), Command::Unknown);
}
#[tokio::test]
async fn resolve_target_uses_env_wave_and_discovery_endpoint() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship", tmp.path(), None);
store.create_wave(&wave).await.expect("seed wave");
std::fs::create_dir_all(tmp.path().join("wave/ship")).expect("wave dir");
std::fs::write(
crate::wave::server::endpoint_path(tmp.path(), "ship"),
"127.0.0.1:4242",
)
.expect("endpoint pointer");
let resolved = resolve_target(
&WaveTargetArgs::default(),
Some(&store),
None,
Some(wave.id().as_str()),
None,
)
.await
.expect("resolve")
.expect("wave context");
assert_eq!(resolved.name, "ship");
assert_eq!(resolved.endpoint.as_deref(), Some("127.0.0.1:4242"));
}
#[tokio::test]
async fn resolve_target_uses_a_hand_set_name_env() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let (addr, _runtime, _inbox) = boot_server(&origin, "ship").await;
let wave = make_wave("ship", &origin, None);
store.create_wave(&wave).await.expect("seed wave");
let resolved = resolve_target(
&WaveTargetArgs::default(),
Some(&store),
None,
Some("ship"),
None,
)
.await
.expect("resolve")
.expect("hand-set name is a wave context");
assert_eq!(resolved.name, "ship");
assert_eq!(resolved.endpoint.as_deref(), Some(addr.as_str()));
}
#[tokio::test]
async fn resolve_target_errors_on_a_stale_wave_id() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let stale = crate::id::WaveId::new().to_string();
let err = resolve_target(
&WaveTargetArgs::default(),
Some(&store),
None,
Some(&stale),
None,
)
.await
.expect_err("stale id is a loud error");
let message = err.to_string();
assert!(message.contains("stale"), "{message}");
assert!(message.contains(&stale), "{message}");
}
#[tokio::test]
async fn no_wave_context_drops_the_message_with_exit_zero() {
let tmp = tempfile::tempdir().expect("tempdir");
let context = CliContext {
store: None,
repo: Some(tmp.path().to_path_buf()),
env_wave_id: None,
env_channel: None,
};
let resolved = resolve_target(
&WaveTargetArgs::default(),
None,
Some(tmp.path()),
None,
None,
)
.await
.expect("resolve");
assert!(resolved.is_none(), "plain temp dir is not a wave context");
run_with_context(
&context,
&["hello".into()],
false,
&WaveTargetArgs::default(),
)
.await
.expect("dropped publish exits 0");
}
#[tokio::test]
async fn steer_flag_requests_live_steering() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let (_addr, runtime, mut inbox) = boot_server(&origin, "ship").await;
let wave = make_wave("ship", &origin, None);
store.create_wave(&wave).await.expect("seed wave");
let context = CliContext {
store: Some(store),
repo: None,
env_wave_id: None,
env_channel: None,
};
run_with_context(
&context,
&["skip".into(), "the".into(), "migration".into()],
true,
&WaveTargetArgs {
wave: Some("ship".into()),
parent: false,
},
)
.await
.expect("post human message");
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].text, "skip the migration");
assert_eq!(thread[0].from, None);
let InboxItem::Message(message) = inbox.try_recv().expect("steer inbox item") else {
panic!("expected message inbox item");
};
assert_eq!(message.op, MessageOp::Steer);
}
#[tokio::test]
async fn plain_chat_is_unattributed_like_the_mac_composer() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let (_addr, runtime, mut inbox) = boot_server(&origin, "ship").await;
let wave = make_wave("ship", &origin, None);
store.create_wave(&wave).await.expect("seed wave");
let context = CliContext {
store: Some(store),
repo: None,
env_wave_id: None,
env_channel: None,
};
run_with_context(
&context,
&["CI".into(), "failed".into()],
false,
&WaveTargetArgs {
wave: Some("ship".into()),
parent: false,
},
)
.await
.expect("post plain message");
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].text, "CI failed");
assert_eq!(thread[0].from, None, "human turns carry no byline");
let InboxItem::Message(message) = inbox.try_recv().expect("inbox item") else {
panic!("expected message inbox item");
};
assert_eq!(message.op, MessageOp::Message);
}
#[tokio::test]
async fn parent_targeting_reaches_the_parent_server_unattributed() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("parent-repo");
std::fs::create_dir_all(&origin).unwrap();
let (addr, parent_runtime, mut parent_inbox) = boot_server(&origin, "goals").await;
let parent = make_wave("goals", &origin, None);
store.create_wave(&parent).await.expect("seed parent");
let child = make_wave("concerto", tmp.path(), Some(parent.id()));
store.create_wave(&child).await.expect("seed child");
let resolved = resolve_target(
&WaveTargetArgs {
wave: None,
parent: true,
},
Some(&store),
None,
Some(child.id().as_str()),
None,
)
.await
.expect("resolve parent")
.expect("wave context");
assert_eq!(resolved.name, "goals");
let endpoint = resolved.require_endpoint().expect("live endpoint");
assert_eq!(endpoint, addr);
let refused = reqwest::Client::new()
.post(format!("http://{endpoint}/messages"))
.json(&serde_json::json!({
"op": "say",
"text": "blocked on the schema",
"from": "wave concerto",
}))
.send()
.await
.expect("post");
assert_eq!(refused.status(), reqwest::StatusCode::BAD_REQUEST);
post_json(
&endpoint,
"/messages",
&serde_json::json!({ "op": "message", "text": "blocked on the schema" }),
)
.await
.expect("post message");
let thread = parent_runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].text, "blocked on the schema");
assert_eq!(thread[0].from, None);
let InboxItem::Message(msg) = parent_inbox.try_recv().expect("inbox item") else {
panic!("expected message inbox item");
};
assert_eq!(msg.op, MessageOp::Message);
let (_, events) = crate::wave::journal::Journal::open(&crate::wave::journal::journal_path(
&origin, "goals",
))
.expect("journal");
assert!(matches!(
&events[0].kind,
EventKind::UserMessage {
op: MessageOp::Message,
from: None,
..
}
));
}
#[tokio::test]
async fn parent_of_a_root_wave_is_a_clear_error() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let root = make_wave("goals", tmp.path(), None);
store.create_wave(&root).await.expect("seed root");
let err = resolve_target(
&WaveTargetArgs {
wave: None,
parent: true,
},
Some(&store),
None,
Some(root.id().as_str()),
None,
)
.await
.expect_err("root has no parent");
assert!(
err.to_string().contains("wave 'goals' has no parent"),
"error names the root wave: {err}"
);
}
#[tokio::test]
async fn a_dotted_name_resolves_to_the_family_head() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let (addr, _runtime, _inbox) = boot_server(&origin, "ship").await;
let wave = make_wave("ship", &origin, None);
store.create_wave(&wave).await.expect("seed wave");
let resolved = resolve_target(
&WaveTargetArgs {
wave: Some("ship.148e".into()),
parent: false,
},
Some(&store),
None,
None,
None,
)
.await
.expect("resolve")
.expect("wave context");
assert_eq!(resolved.name, "ship", "the family head is the wave");
assert_eq!(resolved.endpoint.as_deref(), Some(addr.as_str()));
}
#[tokio::test]
async fn no_live_server_is_a_clear_error() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship", tmp.path(), None);
store.create_wave(&wave).await.expect("seed wave");
let resolved = resolve_target(
&WaveTargetArgs::default(),
Some(&store),
None,
Some(wave.id().as_str()),
None,
)
.await
.expect("resolve")
.expect("wave context");
let err = resolved.require_endpoint().expect_err("no server");
let message = err.to_string();
assert!(message.contains("no live listener"), "{message}");
assert!(message.contains("lf wave ship"), "{message}");
assert!(message.contains("not implemented yet"), "{message}");
}
}