pub(crate) mod channel;
pub mod journal;
pub(crate) mod memory;
pub mod mind;
pub(crate) mod registry;
pub mod resident;
pub mod runtime;
pub mod server;
pub mod state;
pub mod subscription;
pub(crate) mod supervisor;
pub mod wire;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use anyhow::{anyhow, Result};
use crate::engine::repo::find_repo_root;
use crate::engine::worktrees::main_repo_root;
use crate::lfd::types::WAVE_SERVER_ENDPOINT_ENV;
use crate::lfdb::{open_existing_store, SharedStore};
use crate::ops::util::resolve_wave_name;
use crate::wave::runtime::WaveRuntime;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MindPolicy {
Spawn,
Dormant,
}
pub fn run(name: &str, force: bool, no_mind: bool, mind_only: bool) -> Result<()> {
if mind_only {
return resident::run(name);
}
let repo_root = find_repo_root()?;
let main_repo = main_repo_root(&repo_root).unwrap_or_else(|_| repo_root.clone());
let wave = resolve_wave_name(&main_repo, Some(name))
.ok_or_else(|| anyhow!("invalid wave name: '{name}'"))?;
let mind = if no_mind {
MindPolicy::Dormant
} else {
MindPolicy::Spawn
};
let rt = tokio::runtime::Runtime::new()?;
rt.block_on(async {
let registry_config = resolve_registry(&main_repo, &wave, force).await;
serve(
main_repo,
wave,
registry_config,
force,
mind,
shutdown_signal(),
)
.await
})
}
async fn resolve_registry(
main_repo: &Path,
wave: &str,
force: bool,
) -> Option<registry::RegistryConfig> {
let Some(store) = open_existing_store().await else {
tracing::warn!(
wave,
"no session registry on this machine; running unregistered \
(no one-brain enforcement, no worker observations)"
);
return None;
};
let store: SharedStore = Arc::new(store);
match registry::ensure_wave_row(&store, main_repo, wave).await {
Ok(row) => Some(registry::RegistryConfig {
store,
wave: row,
cwd: main_repo.display().to_string(),
pid: std::process::id(),
force,
}),
Err(err) => {
tracing::warn!(wave, error = %err, "session registry unusable; running unregistered");
None
}
}
}
fn resident_spawner(
wave: String,
repo_root: PathBuf,
endpoint: String,
token: String,
session_env: Vec<(String, String)>,
) -> supervisor::SpawnResident {
Box::new(move || {
let exe = std::env::current_exe()?;
let mut command = tokio::process::Command::new(exe);
command
.arg("wave")
.arg(&wave)
.arg("--mind-only")
.current_dir(&repo_root)
.env(WAVE_SERVER_ENDPOINT_ENV, &endpoint)
.env(wire::RESIDENT_TOKEN_ENV, &token)
.env("PATH", mind::path_for_children())
.stdin(std::process::Stdio::null());
for (key, value) in &session_env {
command.env(key, value);
}
command.spawn()
})
}
async fn serve(
repo_root: PathBuf,
wave: String,
registry_config: Option<registry::RegistryConfig>,
force: bool,
mind: MindPolicy,
shutdown: impl Future<Output = ()> + Send + 'static,
) -> Result<()> {
if let Some(live) = server::live_endpoint(&repo_root, &wave).await {
if !force {
return Err(anyhow!(
"refusing to start: wave '{wave}' already has a live server at \
http://{live} (per wave/{wave}/{endpoint}). Stop that session (or let \
it finish), or rerun with --force to take over.",
endpoint = server::ENDPOINT_FILE,
));
}
tracing::warn!(
wave,
live,
"--force: taking over a live wave server endpoint"
);
}
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
let addr = listener.local_addr()?;
let mut registration: Option<registry::Registration> = None;
let mut registered: Option<(SharedStore, crate::lfd::id::LfdId)> = None;
let mut session_env: Vec<(String, String)> = Vec::new();
if let Some(config) = registry_config {
match registry::register(&config, &addr.to_string()).await {
Ok(registry::RegisterOutcome::Registered(reg)) => {
let reg = *reg;
tracing::info!(
wave,
session_id = %reg.session_id(),
"registered in the session registry as the wave's agent session"
);
session_env = vec![
(
crate::lf::session::WAVE_ID_ENV.to_string(),
config.wave.id().to_string(),
),
(
crate::lf::session::SESSION_ID_ENV.to_string(),
reg.session_id().to_string(),
),
(
crate::lf::session::SESSION_INHERITED_ENV.to_string(),
"1".to_string(),
),
];
let hook_registration = reg.clone();
crate::engine::agent::register_interrupt_cleanup(move || {
hook_registration.deregister_blocking();
});
registration = Some(reg);
registered = Some((config.store.clone(), config.wave.id().clone()));
}
Ok(registry::RegisterOutcome::Refused { message }) => {
return Err(anyhow!(
"refusing to start: {message}. Stop that session (or let it \
finish), or rerun with --force to take over."
));
}
Err(err) => {
tracing::warn!(
wave,
error = %err,
"session registry write failed; running unregistered (no one-brain \
enforcement, no worker observations)"
);
}
}
}
let runtime = WaveRuntime::open(wave.clone(), repo_root.clone())?;
runtime.journal_server_started(std::process::id(), &addr.to_string());
let mut observer: Option<Arc<registry::StoreObserver>> = None;
let mut observer_task: Option<tokio::task::JoinHandle<()>> = None;
if let Some((store, wave_id)) = registered {
let obs = Arc::new(registry::StoreObserver::new(
runtime.clone(),
store,
wave_id,
));
observer_task = Some(tokio::spawn(Arc::clone(&obs).run(registry::POLL_CADENCE)));
observer = Some(obs);
}
let token = server::generate_resident_token();
let door = server::ResidentDoor::new(token.clone());
server::write_resident_token(&repo_root, &wave, &token)?;
let spawner = match mind {
MindPolicy::Spawn => Some(resident_spawner(
wave.clone(),
repo_root.clone(),
addr.to_string(),
token.clone(),
session_env,
)),
MindPolicy::Dormant => None,
};
let supervisor = supervisor::Supervisor::new(
runtime.clone(),
door.clone(),
spawner,
supervisor::SupervisorConfig::default(),
);
let supervisor_handle = supervisor.handle();
let supervisor_task = tokio::spawn(supervisor.run());
server::write_endpoint(&repo_root, &wave, addr)?;
let own_addr = addr.to_string();
let cleanup_repo = repo_root.clone();
let cleanup_wave = wave.clone();
let cleanup_addr = own_addr.clone();
let cleanup_token = token.clone();
let cleanup_door = door.clone();
crate::engine::agent::register_interrupt_cleanup(move || {
if let Some(pid) = cleanup_door.seat_pid() {
supervisor::terminate_resident_blocking(pid);
}
server::remove_endpoint(&cleanup_repo, &cleanup_wave, &cleanup_addr);
server::remove_resident_token(&cleanup_repo, &cleanup_wave, &cleanup_token);
});
println!(
"lf wave · {wave} · listener on http://{addr}{} \
(Ctrl-C to stop, RUST_LOG=loopflow=debug for the firehose)",
match mind {
MindPolicy::Spawn => " · spawning resident",
MindPolicy::Dormant => " · dormant (--no-mind)",
}
);
let app = server::router(
runtime.clone(),
door.clone(),
observer,
Some(supervisor_handle),
);
let result = axum::serve(listener, app)
.with_graceful_shutdown(shutdown)
.await;
supervisor_task.abort();
if let Some(pid) = door.seat_pid() {
supervisor::terminate_resident(pid).await;
}
if let Some(task) = observer_task {
task.abort();
}
if let Some(registration) = registration {
registration.deregister().await;
}
server::remove_endpoint(&repo_root, &wave, &own_addr);
server::remove_resident_token(&repo_root, &wave, &token);
result.map_err(|err| anyhow!("wave server error: {err}"))
}
async fn shutdown_signal() {
let _ = tokio::signal::ctrl_c().await;
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use crate::chat::turns::{ChatRole, ChatTurn};
use crate::chat::types::Lifecycle;
use crate::wave::journal::MessageOp;
use crate::wave::server::ResidentDoor;
use crate::wave::wire::{ResidentDelta, RESIDENT_TOKEN_HEADER};
fn progress_turn(text: &str) -> ChatTurn {
ChatTurn {
id: String::new(),
role: ChatRole::Assistant,
text: text.to_string(),
status: Lifecycle::Completed,
items: Vec::new(),
created_at: "1970-01-01T00:00:00Z".to_string(),
from: None,
}
}
fn narrate(runtime: &WaveRuntime, text: &str) {
runtime.append_finalized_turn(progress_turn(text), Vec::new());
}
async fn boot() -> (String, Arc<WaveRuntime>, tempfile::TempDir) {
let tmp = tempfile::tempdir().expect("tempdir");
let dir = tmp.path().join("wave/ship");
std::fs::create_dir_all(&dir).expect("dir");
std::fs::write(dir.join("MEMORY.md"), "Goal: ship the reactive server.\n").expect("mem");
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open");
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(), ResidentDoor::new("test-token"), None, None);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
(format!("http://{addr}"), runtime, tmp)
}
async fn wait_for<F: Fn() -> bool>(cond: F) {
for _ in 0..200 {
if cond() {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("condition not met in time");
}
#[tokio::test]
async fn finalized_turn_appears_in_conversation() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "Implemented the reactive server.");
let body = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.text()
.await
.unwrap();
assert!(body.contains("Implemented the reactive server."));
assert!(body.contains("\"role\":\"assistant\""));
}
#[tokio::test]
async fn conversation_limit_tails_the_thread() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "first");
narrate(&runtime, "second");
narrate(&runtime, "third");
let body: serde_json::Value = reqwest::get(format!("{base}/conversation?limit=2"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns = body["turns"].as_array().unwrap();
assert_eq!(turns.len(), 2);
assert_eq!(turns[0]["text"], "second");
assert_eq!(turns[1]["text"], "third");
for url in [
format!("{base}/conversation?limit=99"),
format!("{base}/conversation"),
] {
let body: serde_json::Value = reqwest::get(url).await.unwrap().json().await.unwrap();
assert_eq!(body["turns"].as_array().unwrap().len(), 3);
}
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "in progress".into(),
});
let body: serde_json::Value = reqwest::get(format!("{base}/conversation?limit=1"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns = body["turns"].as_array().unwrap();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0]["status"], "running");
assert_eq!(turns[0]["text"], "in progress");
}
#[tokio::test]
async fn posted_message_appears_as_user_turn() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let body: serde_json::Value = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "message", "text": "how's it going?" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
let posted: ChatTurn = serde_json::from_value(body["turn"].clone()).unwrap();
assert_eq!(posted.role, ChatRole::User);
assert_eq!(posted.text, "how's it going?");
assert_eq!(body["state"], "idle");
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].role, ChatRole::User);
}
#[tokio::test]
async fn posted_say_is_attributed_on_the_wire() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let body: serde_json::Value = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({
"op": "say",
"text": "run-7 landed: PR #12",
"from": { "session_id": "sess-7", "label": "worker" },
}))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["turn"]["from"], "worker");
assert_eq!(body["turn"]["role"], "user");
let conversation: serde_json::Value = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(conversation["turns"][0]["from"], "worker");
let (_, events) =
journal::Journal::open(&journal::journal_path(runtime.repo_root(), "ship"))
.expect("journal");
assert_eq!(journal::fold_thread(&events).pending_messages.len(), 1);
let missing_from = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "say", "text": "anon" }))
.send()
.await
.unwrap();
assert_eq!(missing_from.status(), reqwest::StatusCode::BAD_REQUEST);
let stray_from = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({
"op": "message",
"text": "hello",
"from": { "session_id": null, "label": "cli" },
}))
.send()
.await
.unwrap();
assert_eq!(stray_from.status(), reqwest::StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn memory_routes_read_and_write_through_the_server() {
let (base, runtime, _tmp) = boot().await;
let body: serde_json::Value = reqwest::get(format!("{base}/memory"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["content"], "Goal: ship the reactive server.\n");
let client = reqwest::Client::new();
let body: serde_json::Value = client
.post(format!("{base}/memory"))
.json(&serde_json::json!({
"op": "update",
"content": "Rewritten.\n",
"summary": null,
}))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["summary"], "Rewritten.");
assert_eq!(runtime.memory().read(), "Rewritten.\n");
let empty = client
.post(format!("{base}/memory"))
.json(&serde_json::json!({ "op": "add", "content": " ", "summary": null }))
.send()
.await
.unwrap();
assert_eq!(empty.status(), reqwest::StatusCode::BAD_REQUEST);
client
.post(format!("{base}/memory"))
.json(&serde_json::json!({ "op": "add", "content": "one fact", "summary": null }))
.send()
.await
.unwrap();
assert_eq!(runtime.memory().read(), "Rewritten.\n- one fact\n");
}
#[tokio::test]
async fn post_without_op_or_without_text_is_rejected() {
let (base, _runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let no_op = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "text": "hello" }))
.send()
.await
.unwrap();
assert_eq!(no_op.status(), reqwest::StatusCode::UNPROCESSABLE_ENTITY);
for op in ["message", "steer"] {
let empty = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": op, "text": " " }))
.send()
.await
.unwrap();
assert_eq!(
empty.status(),
reqwest::StatusCode::BAD_REQUEST,
"empty text rejected for {op}"
);
}
}
#[tokio::test]
async fn bare_interrupt_while_idle_is_a_noop_success_with_state() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let body: serde_json::Value = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "interrupt", "text": "" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert!(body["turn"].is_null(), "nothing said, nothing appended");
assert_eq!(body["state"], "idle");
assert!(runtime.thread_snapshot().is_empty());
}
#[tokio::test]
async fn health_splits_channel_liveness_from_the_mind() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "first");
let body: serde_json::Value = reqwest::get(format!("{base}/health"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["status"], "serving", "status is channel liveness");
assert!(body["mind"].is_null(), "dormant channel has no mind");
assert_eq!(body["wave"], "ship");
assert_eq!(body["turns"], 1);
runtime.set_resident_expected();
let body: serde_json::Value = reqwest::get(format!("{base}/health"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["mind"], "idle", "mind is the resident's state");
runtime.transition(
crate::wave::state::MindState::Failed {
reason: "vendor gone".into(),
},
"test",
);
let body: serde_json::Value = reqwest::get(format!("{base}/health"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["status"], "serving");
assert_eq!(body["mind"], "failed");
}
#[tokio::test]
async fn resident_door_gates_attaches_and_applies_deltas() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
for request in [
client
.post(format!("{base}/resident/attach"))
.json(&serde_json::json!({ "pid": 1234 })),
client
.post(format!("{base}/resident/attach"))
.header(RESIDENT_TOKEN_HEADER, "wrong")
.json(&serde_json::json!({ "pid": 1234 })),
] {
let denied = request.send().await.unwrap();
assert_eq!(denied.status(), reqwest::StatusCode::UNAUTHORIZED);
}
assert!(!runtime.resident_expected());
runtime.set_resident_expected();
runtime.transition(
crate::wave::state::MindState::Failed {
reason: "old resident died".into(),
},
"test",
);
let attach: serde_json::Value = client
.post(format!("{base}/resident/attach"))
.header(RESIDENT_TOKEN_HEADER, "test-token")
.json(&serde_json::json!({ "pid": std::process::id() }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(attach["wave"], "ship");
assert!(attach["thread_id"].is_null());
assert_eq!(runtime.mind_state().name(), "idle", "attach revives");
let deltas: serde_json::Value = client
.post(format!("{base}/resident/deltas"))
.header(RESIDENT_TOKEN_HEADER, "test-token")
.json(&serde_json::json!({ "deltas": [
{ "kind": "thread_started", "vendor": "codex", "thread_id": "t-1" },
{ "kind": "turn_opened", "answers": [] },
{ "kind": "turn_text", "text": "over the wire" },
{ "kind": "turn_usage", "input_tokens": 7, "output_tokens": 3, "cache_read_tokens": null },
{ "kind": "turn_finished", "status": "completed", "cost_usd": 0.01 },
] }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(deltas["accepted"], 5);
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].text, "over the wire");
assert_eq!(thread[0].status, Lifecycle::Completed);
runtime.journal_run_observed("run-1", "sess-1", "implement", "wire it");
let context: serde_json::Value = client
.get(format!("{base}/resident/context"))
.header(RESIDENT_TOKEN_HEADER, "test-token")
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(context["thread_id"], "t-1");
assert_eq!(context["in_flight"][0]["run_id"], "run-1");
assert_eq!(context["in_flight"][0]["flow"], "implement");
}
#[tokio::test]
async fn sse_replays_on_connect_then_streams_live() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "replayed turn");
let host = base.strip_prefix("http://").unwrap().to_string();
let mut stream = tokio::net::TcpStream::connect(&host).await.unwrap();
stream
.write_all(b"GET /events HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n")
.await
.unwrap();
narrate(&runtime, "live turn");
let mut acc = String::new();
let mut buf = [0u8; 2048];
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
let read = tokio::time::timeout_at(deadline, stream.read(&mut buf)).await;
match read {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(n)) => {
acc.push_str(&String::from_utf8_lossy(&buf[..n]));
if acc.contains("replayed turn") && acc.contains("live turn") {
break;
}
}
Ok(Err(_)) => break,
}
}
assert!(acc.contains("event: turn"), "SSE frames are named `turn`");
assert!(
acc.contains("replayed turn"),
"replays the thread on connect"
);
assert!(
acc.contains("live turn"),
"streams turns narrated after connect"
);
assert!(
!acc.contains("event: inbox"),
"the default stream carries no inbox frames"
);
}
#[tokio::test]
async fn events_inbox_scope_replays_pending_and_streams_ops() {
let (base, runtime, _tmp) = boot().await;
runtime.deliver_user_message("queued before".into(), MessageOp::Message);
let host = base.strip_prefix("http://").unwrap().to_string();
let mut stream = tokio::net::TcpStream::connect(&host).await.unwrap();
stream
.write_all(
b"GET /events?inbox=true HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n",
)
.await
.unwrap();
let mut acc = String::new();
let mut buf = [0u8; 4096];
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut interrupt_sent = false;
loop {
if acc.contains("queued before") && !interrupt_sent {
interrupt_sent = true;
runtime.deliver_interrupt();
}
let read = tokio::time::timeout_at(deadline, stream.read(&mut buf)).await;
match read {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(n)) => {
acc.push_str(&String::from_utf8_lossy(&buf[..n]));
if acc.contains("queued before") && acc.contains(r#""id":null"#) {
break;
}
}
Ok(Err(_)) => break,
}
}
assert!(acc.contains("event: inbox"), "inbox frames are named");
assert!(
acc.contains("queued before") && acc.contains(r#""op":"message""#),
"the pending queue replays: {acc}"
);
assert!(
acc.contains(r#""id":null"#) && acc.contains(r#""op":"interrupt""#),
"a live bare interrupt rides id-less: {acc}"
);
}
struct SseClient {
stream: tokio::net::TcpStream,
raw: Vec<u8>,
}
impl SseClient {
async fn connect(base: &str) -> Self {
let host = base.strip_prefix("http://").unwrap();
let mut stream = tokio::net::TcpStream::connect(host).await.unwrap();
stream
.write_all(
b"GET /events HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n",
)
.await
.unwrap();
Self {
stream,
raw: Vec::new(),
}
}
async fn frames_until(&mut self, pred: impl Fn(&[ChatTurn]) -> bool) -> Vec<ChatTurn> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut buf = [0u8; 4096];
loop {
let frames = parse_turn_frames(&dechunk(&self.raw));
if pred(&frames) {
return frames;
}
match tokio::time::timeout_at(deadline, self.stream.read(&mut buf)).await {
Ok(Ok(0)) | Err(_) => {
panic!("SSE ended before condition; {} frames so far", frames.len())
}
Ok(Ok(n)) => self.raw.extend_from_slice(&buf[..n]),
Ok(Err(err)) => panic!("SSE read error: {err}"),
}
}
}
async fn states_until(&mut self, pred: impl Fn(&[String]) -> bool) -> Vec<String> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut buf = [0u8; 4096];
loop {
let states = parse_state_frames(&dechunk(&self.raw));
if pred(&states) {
return states;
}
match tokio::time::timeout_at(deadline, self.stream.read(&mut buf)).await {
Ok(Ok(0)) | Err(_) => {
panic!("SSE ended before condition; states so far: {states:?}")
}
Ok(Ok(n)) => self.raw.extend_from_slice(&buf[..n]),
Ok(Err(err)) => panic!("SSE read error: {err}"),
}
}
}
}
fn dechunk(raw: &[u8]) -> String {
let text = String::from_utf8_lossy(raw);
let Some(head_end) = text.find("\r\n\r\n") else {
return String::new();
};
let mut body = &text[head_end + 4..];
let mut out = String::new();
while let Some(size_end) = body.find("\r\n") {
let Ok(size) = usize::from_str_radix(body[..size_end].trim(), 16) else {
break;
};
let start = size_end + 2;
if size == 0 || body.len() < start + size {
break;
}
out.push_str(&body[start..start + size]);
body = &body[(start + size + 2).min(body.len())..];
}
out
}
fn parse_turn_frames(sse_body: &str) -> Vec<ChatTurn> {
sse_body
.lines()
.filter_map(|line| line.strip_prefix("data:"))
.filter_map(|data| serde_json::from_str(data.trim()).ok())
.collect()
}
fn parse_state_frames(sse_body: &str) -> Vec<String> {
let mut states = Vec::new();
let mut in_state_event = false;
for line in sse_body.lines() {
if let Some(name) = line.strip_prefix("event:") {
in_state_event = name.trim() == "state";
} else if let Some(data) = line.strip_prefix("data:") {
if in_state_event {
states.push(data.trim().to_string());
}
}
}
states
}
#[tokio::test]
async fn sse_late_subscriber_watches_the_open_turn_grow_and_finalize() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "already finalized");
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "thinking".into(),
});
let mut client = SseClient::connect(&base).await;
let frames = client
.frames_until(|f| f.iter().any(|t| t.status == Lifecycle::Running))
.await;
assert!(
frames
.iter()
.any(|t| t.text == "already finalized" && t.status == Lifecycle::Completed),
"replay carries the finalized thread"
);
let open = frames
.iter()
.find(|t| t.status == Lifecycle::Running && t.text == "thinking")
.unwrap()
.clone();
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "more".into(),
});
client
.frames_until(|f| {
f.iter().any(|t| {
t.id == open.id && t.text == "thinking\nmore" && t.status == Lifecycle::Running
})
})
.await;
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
});
let frames = client
.frames_until(|f| {
f.iter()
.any(|t| t.id == open.id && t.status == Lifecycle::Completed)
})
.await;
let last = frames.iter().rfind(|t| t.id == open.id).unwrap();
assert_eq!(last.status, Lifecycle::Completed, "terminal frame is last");
assert_eq!(last.text, "thinking\nmore");
}
#[tokio::test]
async fn sse_state_events_track_the_mind_live() {
let (base, runtime, _tmp) = boot().await;
let mut client = SseClient::connect(&base).await;
let states = client.states_until(|s| !s.is_empty()).await;
assert_eq!(states, vec!["idle"]);
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
let states = client.states_until(|s| s.len() >= 2).await;
assert_eq!(states, vec!["idle", "turning"]);
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
});
let states = client.states_until(|s| s.len() >= 3).await;
assert_eq!(states, vec!["idle", "turning", "idle"]);
}
#[tokio::test]
async fn conversation_includes_the_open_running_turn() {
let (base, runtime, _tmp) = boot().await;
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "half a thought".into(),
});
let body: serde_json::Value = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns = body["turns"].as_array().unwrap();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0]["status"], "running");
assert_eq!(turns[0]["text"], "half a thought");
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
});
let body: serde_json::Value = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns = body["turns"].as_array().unwrap();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0]["status"], "completed");
}
#[tokio::test]
async fn restart_mid_turn_never_serves_a_stale_running_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
{
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open");
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "half a thought".into(),
});
}
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("reopen");
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(), ResidentDoor::new("test-token"), None, None);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
let base = format!("http://{addr}");
let body: serde_json::Value = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns = body["turns"].as_array().unwrap();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0]["status"], "failed", "janitor closed the turn");
assert_eq!(turns[0]["text"], "half a thought");
let mut client = SseClient::connect(&base).await;
let frames = client
.frames_until(|f| f.iter().any(|t| t.status == Lifecycle::Failed))
.await;
assert!(
frames.iter().all(|t| t.status != Lifecycle::Running),
"no stale running turn in replay"
);
}
async fn boot_family(children: &[&str]) -> (String, Arc<WaveRuntime>, tempfile::TempDir) {
let tmp = tempfile::tempdir().expect("tempdir");
let origin = tmp.path().join("repo");
std::fs::create_dir_all(origin.join("wave/ship")).expect("wave dir");
for child in children {
std::fs::create_dir_all(crate::wave::channel::child_worktree_path(&origin, child))
.expect("child worktree");
}
let runtime = WaveRuntime::open("ship".into(), origin).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(), ResidentDoor::new("test-token"), None, None);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
(format!("http://{addr}"), runtime, tmp)
}
#[tokio::test]
async fn posted_message_with_channel_lands_in_that_journal() {
let (base, runtime, _tmp) = boot_family(&["ship.148e"]).await;
let client = reqwest::Client::new();
let body: serde_json::Value = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({
"op": "say",
"text": "landed the parser",
"from": { "session_id": null, "label": "worker" },
"channel": "ship.148e",
}))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["turn"]["text"], "landed the parser");
let child_journal =
crate::wave::channel::child_journal_path(runtime.repo_root(), "ship.148e");
let events = journal::read_events(&child_journal);
assert_eq!(events.len(), 1);
let forwarded = runtime.thread_snapshot();
assert_eq!(forwarded.len(), 1, "the report reached the wave thread");
assert_eq!(forwarded[0].text, "landed the parser");
assert_eq!(
journal::read_events(&journal::journal_path(runtime.repo_root(), "ship"))
.iter()
.filter(|e| matches!(e.kind, journal::EventKind::UserMessage { .. }))
.count(),
1,
"the forwarded report is journaled on the wave channel",
);
client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "message", "text": "to the wave", "channel": "ship" }))
.send()
.await
.unwrap();
assert_eq!(runtime.thread_snapshot().len(), 2);
for channel in ["concerto", "ship.ghost"] {
let refused = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "message", "text": "?", "channel": channel }))
.send()
.await
.unwrap();
assert_eq!(
refused.status(),
reqwest::StatusCode::NOT_FOUND,
"channel '{channel}' bounces"
);
}
}
#[tokio::test]
async fn channels_door_journals_the_opening_once() {
let (base, runtime, _tmp) = boot_family(&[]).await;
let client = reqwest::Client::new();
let body: serde_json::Value = client
.post(format!("{base}/channels"))
.json(&serde_json::json!({ "name": "ship.148e", "run_id": "run-1" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["turn"]["text"], "work line ship.148e opened");
assert_eq!(body["turn"]["from"], "dispatch");
let again: serde_json::Value = client
.post(format!("{base}/channels"))
.json(&serde_json::json!({ "name": "ship.148e", "run_id": "run-1" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert!(again["turn"].is_null(), "repeated knock appends nothing");
assert_eq!(runtime.thread_snapshot().len(), 1);
let foreign = client
.post(format!("{base}/channels"))
.json(&serde_json::json!({ "name": "concerto.x", "run_id": "run-2" }))
.send()
.await
.unwrap();
assert_eq!(foreign.status(), reqwest::StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn events_family_subscription_carries_channel_tagged_frames() {
let (base, runtime, _tmp) = boot_family(&["ship.a", "ship.b"]).await;
narrate(&runtime, "wave turn");
runtime
.deliver_to_channel(
"ship.a",
journal::MessageOp::Message,
"a replay".into(),
None,
)
.unwrap();
let host = base.strip_prefix("http://").unwrap().to_string();
let mut stream = tokio::net::TcpStream::connect(&host).await.unwrap();
stream
.write_all(b"GET /events HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n")
.await
.unwrap();
runtime
.deliver_to_channel("ship.b", journal::MessageOp::Message, "b live".into(), None)
.unwrap();
let mut acc = String::new();
let mut buf = [0u8; 4096];
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
let read = tokio::time::timeout_at(deadline, stream.read(&mut buf)).await;
match read {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(n)) => {
acc.push_str(&String::from_utf8_lossy(&buf[..n]));
if acc.contains("wave turn")
&& acc.contains("a replay")
&& acc.contains("b live")
{
break;
}
}
Ok(Err(_)) => break,
}
}
assert!(acc.contains("a replay"), "child replay arrives: {acc}");
assert!(acc.contains("b live"), "child live frame arrives");
assert!(
acc.contains(r#""channel":"ship.a""#) && acc.contains(r#""channel":"ship.b""#),
"child frames carry their channel tag"
);
let wave_frame = acc
.lines()
.find(|line| line.contains("wave turn"))
.expect("wave frame");
assert!(
!wave_frame.contains(r#""channel""#),
"primary frames are untagged: {wave_frame}"
);
let mut one = tokio::net::TcpStream::connect(&host).await.unwrap();
one.write_all(
b"GET /events?channel=ship.a HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n",
)
.await
.unwrap();
let mut acc_one = String::new();
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
let read = tokio::time::timeout_at(deadline, one.read(&mut buf)).await;
match read {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(n)) => {
acc_one.push_str(&String::from_utf8_lossy(&buf[..n]));
if acc_one.contains("a replay") {
break;
}
}
Ok(Err(_)) => break,
}
}
assert!(acc_one.contains("a replay"));
assert!(!acc_one.contains("wave turn"), "no primary frames");
assert!(
!acc_one.contains("event: state"),
"no mind on a child channel"
);
let refused = reqwest::get(format!("{base}/events?channel=concerto"))
.await
.unwrap();
assert_eq!(refused.status(), reqwest::StatusCode::NOT_FOUND);
let both = reqwest::get(format!("{base}/events?channel=ship&prefix=ship"))
.await
.unwrap();
assert_eq!(both.status(), reqwest::StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn serve_publishes_and_removes_discovery_pointer_and_token() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let handle = tokio::spawn(async move {
serve(
repo2,
"ship".into(),
None,
false,
MindPolicy::Dormant,
async {
let _ = shutdown_rx.await;
},
)
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let contents = std::fs::read_to_string(&endpoint).unwrap();
assert!(
contents.starts_with("127.0.0.1:"),
"pointer is just an address"
);
assert!(
server::read_resident_token(&repo, "ship").is_some(),
"the resident token publishes beside the pointer"
);
shutdown_tx.send(()).unwrap();
handle.await.unwrap().unwrap();
assert!(!endpoint.exists(), "pointer removed on shutdown");
assert!(
server::read_resident_token(&repo, "ship").is_none(),
"token removed on shutdown"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn resolve_registry_without_a_store_runs_unregistered() {
let _env = crate::lf::session::test_env_lock();
let tmp = tempfile::tempdir().expect("tempdir");
let previous = std::env::var_os("LFD_DB_PATH");
std::env::set_var("LFD_DB_PATH", tmp.path().join("absent.db"));
let config = resolve_registry(tmp.path(), "ship", false).await;
match previous {
Some(value) => std::env::set_var("LFD_DB_PATH", value),
None => std::env::remove_var("LFD_DB_PATH"),
}
assert!(config.is_none(), "missing store boots unregistered");
}
#[tokio::test]
async fn serve_overwrites_a_stale_endpoint_file() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let dead = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let dead_addr = dead.local_addr().unwrap();
drop(dead);
server::write_endpoint(&repo, "ship", dead_addr).expect("stale pointer");
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let handle = tokio::spawn(async move {
serve(
repo2,
"ship".into(),
None,
false,
MindPolicy::Dormant,
async {
let _ = shutdown_rx.await;
},
)
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| {
std::fs::read_to_string(&endpoint)
.is_ok_and(|contents| contents != dead_addr.to_string())
})
.await;
shutdown_tx.send(()).unwrap();
handle.await.unwrap().unwrap();
}
#[tokio::test]
async fn serve_refuses_to_start_over_a_live_endpoint_file() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let first = tokio::spawn(async move {
serve(
repo2,
"ship".into(),
None,
false,
MindPolicy::Dormant,
async {
let _ = shutdown_rx.await;
},
)
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let first_addr = std::fs::read_to_string(&endpoint).unwrap();
let (_shutdown_tx2, shutdown_rx2) = tokio::sync::oneshot::channel::<()>();
let err = serve(
repo.clone(),
"ship".into(),
None,
false,
MindPolicy::Dormant,
async {
let _ = shutdown_rx2.await;
},
)
.await
.expect_err("live endpoint refuses a second server");
assert!(
err.to_string().contains("--force"),
"error points at --force: {err}"
);
assert_eq!(
std::fs::read_to_string(&endpoint).unwrap(),
first_addr,
"refused boot never touched the pointer"
);
shutdown_tx.send(()).unwrap();
first.await.unwrap().unwrap();
assert!(!endpoint.exists(), "first server still owns its shutdown");
}
fn make_wave_row(name: &str) -> crate::lfd::types::Wave {
use crate::lfd::types::{RepoWork, WaveStatus};
crate::lfd::types::Wave {
id: crate::lfd::id::LfdId::new(),
name: name.to_string(),
primary_flow: "ship-roadmap".to_string(),
goal: "ship-roadmap".to_string(),
metrics: Vec::new(),
repos: vec![RepoWork {
repo: "/tmp/repo".to_string(),
worktree: String::new(),
branch: String::new(),
status: WaveStatus::Idle,
iteration: 0,
cycle_start_iteration: 0,
position: 0,
}],
direction: Vec::new(),
area: Vec::new(),
paused: false,
created_at: Some(time::OffsetDateTime::now_utc()),
workers: 1,
parent_wave_id: None,
}
}
#[tokio::test]
async fn serve_registers_the_brain_and_deregisters_on_shutdown() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let store: crate::lfdb::SharedStore = Arc::new(
crate::lfdb::open_store(&crate::lfdb::StorageConfig::sqlite(
tmp.path().join("lfd.db"),
))
.await
.expect("open store"),
);
let wave_row = make_wave_row("ship");
store.create_wave(&wave_row).await.expect("seed wave");
let wave_id = wave_row.id().clone();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let config = registry::RegistryConfig {
store: store.clone(),
wave: wave_row,
cwd: repo.display().to_string(),
pid: std::process::id(),
force: false,
};
let handle = tokio::spawn(async move {
serve(
repo2,
"ship".into(),
Some(config),
false,
MindPolicy::Dormant,
async {
let _ = shutdown_rx.await;
},
)
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let live = store
.live_wave_agent_session(&wave_id)
.await
.expect("live lookup")
.expect("brain registered");
let addr = std::fs::read_to_string(&endpoint).unwrap();
assert_eq!(
live.env
.get(crate::lfd::types::WAVE_SERVER_ENDPOINT_ENV)
.map(String::as_str),
Some(addr.trim())
);
shutdown_tx.send(()).unwrap();
handle.await.unwrap().unwrap();
assert!(!endpoint.exists(), "pointer removed on shutdown");
assert!(
store
.live_wave_agent_session(&wave_id)
.await
.expect("live lookup")
.is_none(),
"graceful shutdown deregisters the brain"
);
}
#[tokio::test]
async fn serve_refuses_to_start_over_a_live_brain() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let store: crate::lfdb::SharedStore = Arc::new(
crate::lfdb::open_store(&crate::lfdb::StorageConfig::sqlite(
tmp.path().join("lfd.db"),
))
.await
.expect("open store"),
);
let wave_row = make_wave_row("ship");
store.create_wave(&wave_row).await.expect("seed wave");
let first = registry::RegistryConfig {
store: store.clone(),
wave: wave_row.clone(),
cwd: "/tmp/repo".to_string(),
pid: std::process::id(),
force: false,
};
let registry::RegisterOutcome::Registered(_live) =
registry::register(&first, "127.0.0.1:9")
.await
.expect("register")
else {
panic!("first registration succeeds");
};
let (_shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let err = serve(
tmp.path().to_path_buf(),
"ship".into(),
Some(registry::RegistryConfig {
store,
wave: wave_row,
cwd: "/tmp/repo".to_string(),
pid: std::process::id(),
force: false,
}),
false,
MindPolicy::Dormant,
async {
let _ = shutdown_rx.await;
},
)
.await
.expect_err("second brain refused");
assert!(
err.to_string().contains("--force"),
"error points at --force: {err}"
);
}
}