use std::collections::{HashMap, HashSet};
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::time::Duration;
use anyhow::{anyhow, Result};
use chrono::{DateTime, Utc};
use tokio::sync::mpsc;
use tokio::time::Instant;
use crate::chat::types::{ConversationEvent, Lifecycle, TurnUsage};
use crate::engine::flow::{available_flow_names, load_goal, render_goal, GoalRenderContext};
use crate::engine::wave_config::{read_wave_config, WaveCronDef};
use crate::harness::{default_create_harness, ApprovalPolicy, Harness};
use crate::wave::journal::{MessageId, MessageOp, PendingMessage};
use crate::wave::playhead::{BodyProvenance, StepKind, StepOutcome, StepRef};
use crate::wave::resident::ListenerClient;
use crate::wave::runtime::InboxItem;
use crate::wave::supervisor::sleep_until_opt;
use crate::wave::wire::{ProviderSessionRef, ResidentDelta, ResidentStateTo};
pub const HEARTBEAT_IDLE: Duration = Duration::from_secs(4 * 60 * 60);
pub const MAX_CONSECUTIVE_PASS_FAILURES: u32 = 3;
pub const CRON_GRACE: chrono::Duration = chrono::Duration::hours(24);
const HEARTBEAT_PROMPT: &str = "Heartbeat: re-read your goal and memory, then take the next \
orchestration skill. If nothing needs doing, say so in one line.";
fn finish_capture(capture: Option<&crate::trace::CaptureHandle>, outcome: &str) {
let Some(capture) = capture else {
return;
};
if let Err(error) = capture.finish(outcome, false) {
tracing::warn!(%error, %outcome, "failed to finalize trace capture");
}
}
fn read_crons(origin_repo: &Path, wave: &str) -> Vec<WaveCronDef> {
read_wave_config(origin_repo, wave)
.and_then(|config| config.crons)
.unwrap_or_default()
}
fn cron_key(cron: &WaveCronDef) -> String {
format!("{} {}", cron.schedule, cron.flow)
}
fn next_cron_fire(
schedule: &str,
last_fired: Option<DateTime<Utc>>,
now: DateTime<Utc>,
) -> Option<DateTime<Utc>> {
let schedule = cron::Schedule::from_str(schedule).ok()?;
let check_from = last_fired.unwrap_or(now - CRON_GRACE);
schedule.after(&check_from).next()
}
pub(crate) fn cron_prompt(due: &[WaveCronDef]) -> String {
due.iter()
.map(|cron| format!("cron due: {} — dispatch it", cron.flow))
.collect::<Vec<_>>()
.join("\n")
}
#[derive(Debug, Clone)]
pub struct LoopConfig {
pub heartbeat_idle: Duration,
pub pass_timeout: Duration,
pub max_turns: Option<u32>,
}
impl Default for LoopConfig {
fn default() -> Self {
Self {
heartbeat_idle: HEARTBEAT_IDLE,
pass_timeout: Duration::from_secs(30 * 60),
max_turns: Some(20),
}
}
}
pub fn path_for_children() -> OsString {
let inherited = std::env::var_os("PATH").unwrap_or_default();
let Some(exe_dir) = std::env::current_exe()
.ok()
.and_then(|exe| exe.parent().map(Path::to_path_buf))
else {
return inherited;
};
let paths = std::iter::once(exe_dir).chain(std::env::split_paths(&inherited));
std::env::join_paths(paths).unwrap_or(inherited)
}
fn wave_pass_seed(origin_repo: &Path, wave: &str, wake: &str) -> String {
let seed = build_goal_seed(origin_repo, wave);
format!(
"{seed}\n\n{}\n\n{}\n\n<wake>\n{wake}\n</wake>",
orchestration_discipline(wave),
crate::engine::prompt::loopflow_section()
)
}
fn build_goal_seed(repo: &Path, wave: &str) -> String {
let memory = crate::engine::wave_context::gather_wave_memory(repo, wave).unwrap_or_default();
match load_goal(wave, repo) {
Ok(goal) => {
let ctx = GoalRenderContext {
flows: available_flow_names(repo),
memory,
};
render_goal(&goal, &ctx)
}
Err(_) => {
let mem_block = if memory.trim().is_empty() {
"(memory is empty)".to_string()
} else {
memory
};
format!(
"You are the agent of the '{wave}' wave. Drive the wave's goal \
forward.\n\nCurrent memory:\n{mem_block}"
)
}
}
}
fn orchestration_discipline(wave: &str) -> String {
format!(
"You are the loop of the '{wave}' wave — its long-running orchestrator.\n\
Discipline:\n\
- Read state and filed tasks when available. If PM fails, report it \
once and continue from GOAL, MEMORY, and project KRs; do not turn \
infrastructure repair into this wave's work.\n\
- Execute the next move inline by default. Resolve the one concrete \
blocker between the wave and progress in this process.\n\
- Create or select a Linear Project and Linear task before delegating \
file-writing work. Start it with `lf task run <issue-id>`. Never \
delegate anonymous work or the whole wave objective.\n\
- Supervise durable Task Sessions with `lf task status`, `follow-up`, \
`steer`, `interrupt`, `wait`, and `resume`. Each task owns one immutable \
worktree and one PR to main; keep the Wave home free of shipping edits.\n\
- Keep turns centered on selection, direct progress, sequencing, and \
authored reports.\n\
- Trust worker summaries; never re-read worker transcripts.\n\
- A human message is steering: answer it directly and adjust course \
before returning to the goal."
)
}
fn lf_command() -> std::process::Command {
if let Ok(path) = std::env::current_exe() {
return std::process::Command::new(path);
}
std::process::Command::new("lf")
}
fn body_provenance(step: &StepRef, cwd: &Path) -> BodyProvenance {
let configured = crate::engine::load_config_or_default(Some(cwd));
let agent = configured.agent.as_deref().unwrap_or("codex");
let (harness, model) = crate::engine::parse_agent(agent);
let mut body = BodyProvenance::for_step(step, cwd);
body.harness = Some(harness);
body.model = model;
body
}
fn provider_session_id_for_harness(
provider_session: Option<&ProviderSessionRef>,
harness: &str,
) -> Option<String> {
provider_session
.filter(|session| session.harness == harness)
.map(|session| session.session_id.clone())
}
enum LoopEnd {
ListenerGone,
Failed(String),
}
enum InboxAction {
Interrupt { skip: bool },
Deliver(InboxItem),
ListenerGone,
}
enum TimeoutAction {
End,
Expire,
}
#[cfg(test)]
type SpawnPass = Box<
dyn Fn(&Path, &StepRef, &str, Option<u32>) -> std::io::Result<tokio::process::Child> + Send,
>;
type CreateBodyHarness = Box<
dyn Fn(
&str,
ApprovalPolicy,
mpsc::UnboundedSender<ConversationEvent>,
) -> Result<Box<dyn Harness>>
+ Send,
>;
type PrepareBodyHarness = Box<
dyn Fn(&str, &str, &str, Option<u32>) -> Result<crate::lf::commands::run::PreparedHarnessTurn>
+ Send,
>;
enum BodyBackend {
Harness {
prepare: PrepareBodyHarness,
create: CreateBodyHarness,
},
#[cfg(test)]
Process(SpawnPass),
}
pub async fn run_loop(
client: ListenerClient,
inbox_rx: mpsc::UnboundedReceiver<InboxItem>,
cwd: PathBuf,
origin_repo: PathBuf,
wave: String,
config: LoopConfig,
) -> Result<()> {
let backend = BodyBackend::Harness {
prepare: Box::new(crate::lf::commands::run::prepare_harness_turn),
create: Box::new(default_create_harness),
};
run_loop_with(client, inbox_rx, cwd, origin_repo, wave, config, backend).await
}
#[allow(clippy::too_many_arguments)]
async fn run_loop_with(
client: ListenerClient,
mut inbox_rx: mpsc::UnboundedReceiver<InboxItem>,
cwd: PathBuf,
origin_repo: PathBuf,
wave: String,
config: LoopConfig,
backend: BodyBackend,
) -> Result<()> {
let mut wave_loop = WaveLoop {
client,
cwd,
origin_repo,
wave,
config,
queue: Vec::new(),
seen: HashSet::new(),
backend,
consecutive_failures: 0,
idle_since: Instant::now(),
cron_last_fired: HashMap::new(),
provider_session: None,
end: None,
};
while wave_loop.end.is_none() {
if !wave_loop.queue.is_empty() {
wave_loop.start_queued_pass(&mut inbox_rx).await;
continue;
}
let heartbeat_at = wave_loop.heartbeat_deadline();
let cron_at = wave_loop.cron_deadline();
tokio::select! {
biased;
item = inbox_rx.recv() => {
match item {
Some(item) => wave_loop.on_inbox(item).await,
None => wave_loop.end = Some(LoopEnd::ListenerGone),
}
}
_ = sleep_until_opt(cron_at), if cron_at.is_some() => {
wave_loop.on_cron(&mut inbox_rx).await;
}
_ = sleep_until_opt(Some(heartbeat_at)) => {
wave_loop.on_heartbeat(&mut inbox_rx).await;
}
}
}
match wave_loop.end {
Some(LoopEnd::Failed(reason)) => Err(anyhow!(reason)),
_ => Ok(()),
}
}
struct WaveLoop {
client: ListenerClient,
cwd: PathBuf,
origin_repo: PathBuf,
wave: String,
config: LoopConfig,
backend: BodyBackend,
queue: Vec<PendingMessage>,
seen: HashSet<MessageId>,
consecutive_failures: u32,
idle_since: Instant,
cron_last_fired: HashMap<String, DateTime<Utc>>,
provider_session: Option<ProviderSessionRef>,
end: Option<LoopEnd>,
}
impl WaveLoop {
fn heartbeat_deadline(&self) -> Instant {
self.idle_since + self.config.heartbeat_idle
}
fn cron_deadline(&self) -> Option<Instant> {
let now = Utc::now();
let next = read_crons(&self.origin_repo, &self.wave)
.iter()
.filter_map(|cron| {
next_cron_fire(
&cron.schedule,
self.cron_last_fired.get(&cron_key(cron)).copied(),
now,
)
})
.min()?;
let wait = (next - now).to_std().unwrap_or(Duration::ZERO);
Some(Instant::now() + wait)
}
async fn on_inbox(&mut self, item: InboxItem) {
match item {
InboxItem::Message(message) => {
if self.seen.insert(message.id.clone()) {
self.queue.push(message);
}
}
InboxItem::Task(observation) => {
let message = crate::wave::journal::task_observation_message(&observation);
if self.seen.insert(message.id.clone()) {
self.queue.push(message);
}
}
InboxItem::Project(observation) => {
let message = crate::wave::journal::project_observation_message(&observation);
if self.seen.insert(message.id.clone()) {
self.queue.push(message);
}
}
InboxItem::Interrupt | InboxItem::Skip => {}
}
}
async fn on_heartbeat(&mut self, inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>) {
self.run_pass(HEARTBEAT_PROMPT.to_string(), Vec::new(), inbox_rx)
.await;
}
async fn on_cron(&mut self, inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>) {
let now = Utc::now();
let due: Vec<WaveCronDef> = read_crons(&self.origin_repo, &self.wave)
.into_iter()
.filter(|cron| {
next_cron_fire(
&cron.schedule,
self.cron_last_fired.get(&cron_key(cron)).copied(),
now,
)
.is_some_and(|fire_at| fire_at <= now)
})
.collect();
if due.is_empty() {
return;
}
for cron in &due {
self.cron_last_fired.insert(cron_key(cron), now);
}
let prompt = cron_prompt(&due);
self.run_pass(prompt, Vec::new(), inbox_rx).await;
}
async fn start_queued_pass(&mut self, inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>) {
let messages = std::mem::take(&mut self.queue);
let answers: Vec<MessageId> = messages.iter().map(|m| m.id.clone()).collect();
let content = messages
.iter()
.map(|m| match &m.from {
Some(from) => format!("[{from}] {}", m.text),
None => m.text.clone(),
})
.collect::<Vec<_>>()
.join("\n\n");
self.run_pass(content, answers, inbox_rx).await;
}
async fn run_pass(
&mut self,
wake: String,
answers: Vec<MessageId>,
inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>,
) {
let mut answers = answers.into_iter().map(|id| id.0).collect::<Vec<_>>();
let mut invocation: Option<(String, u32)> = None;
loop {
let context = match self.fetch_context().await {
Some(context) => context,
None => return,
};
let Some(step) = context.playhead.now else {
self.fail("playhead has no current step").await;
return;
};
let key = (step.invocation_id.clone(), step.iteration);
if invocation.as_ref().is_some_and(|expected| expected != &key) {
return;
}
invocation.get_or_insert(key);
let completed_index = step.index;
let seed = wave_pass_seed(&self.origin_repo, &self.wave, &wake);
let live_skill = step.kind == StepKind::Skill
&& matches!(&self.backend, BodyBackend::Harness { .. });
if live_skill {
self.run_harness_pass(step, seed, answers, inbox_rx).await;
} else {
self.run_process_pass(step, seed, answers, inbox_rx).await;
}
if self.end.is_some() {
return;
}
let next = match self.fetch_context().await {
Some(context) => context.playhead.now,
None => return,
};
let Some(next) = next else { return };
let same_iteration = invocation.as_ref().is_some_and(|(id, iteration)| {
id == &next.invocation_id && *iteration == next.iteration
});
if !same_iteration || next.index == completed_index {
return;
}
answers = Vec::new();
}
}
async fn run_process_pass(
&mut self,
step: StepRef,
seed: String,
answers: Vec<String>,
inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>,
) {
let body = body_provenance(&step, &self.cwd);
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
if self.end.is_some() {
return;
}
let child = match &self.backend {
BodyBackend::Harness { .. } => {
spawn_wave_step(&self.cwd, &step, &seed, self.config.max_turns)
}
#[cfg(test)]
BodyBackend::Process(spawn) => spawn(&self.cwd, &step, &seed, self.config.max_turns),
};
let child = match child {
Ok(child) => child,
Err(err) => {
self.finish_failed_pass(
&body_id,
&format!("failed to spawn {} / {}: {err:#}", step.flow, step.step),
)
.await;
return;
}
};
let mut wait_task = tokio::spawn(async move { child.wait_with_output().await });
let mut timeout = Box::pin(tokio::time::sleep(self.config.pass_timeout));
loop {
tokio::select! {
biased;
item = inbox_rx.recv() => {
match self.inbox_action(item) {
InboxAction::Interrupt { skip } => {
self.interrupt_child(&body_id, &mut wait_task, skip).await;
return;
}
InboxAction::Deliver(item) => self.on_inbox(item).await,
InboxAction::ListenerGone => {
wait_task.abort();
return;
}
}
}
_ = &mut timeout => {
match self.timeout_action() {
TimeoutAction::End => {
wait_task.abort();
return;
}
TimeoutAction::Expire => {
wait_task.abort();
self.finish_timed_out_pass(&body_id).await;
return;
}
}
}
result = &mut wait_task => {
match result {
Ok(output) => self.on_pass_output(&body_id, output).await,
Err(err) => {
self.finish_failed_pass(
&body_id,
&format!("wave wait task failed: {err:#}"),
)
.await;
}
}
return;
}
}
}
}
async fn run_harness_pass(
&mut self,
step: StepRef,
seed: String,
answers: Vec<String>,
inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>,
) {
let mut body = body_provenance(&step, &self.cwd);
let prepared = match &self.backend {
BodyBackend::Harness { prepare, .. } => {
prepare(&step.step, &seed, &self.wave, self.config.max_turns)
}
#[cfg(test)]
BodyBackend::Process(_) => unreachable!("live skill requires a harness backend"),
};
let prepared = match prepared {
Ok(prepared) => prepared,
Err(err) => {
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
self.finish_failed_pass(
&body_id,
&format!("failed to prepare {} / {}: {err:#}", step.flow, step.step),
)
.await;
return;
}
};
body.harness = Some(prepared.harness.clone());
body.model = prepared.model.clone();
let capture = match crate::journal::trace_capture_context(
&self.cwd,
Some(step.flow.clone()),
Some(step.step.clone()),
) {
Some(context) => match crate::trace::CaptureHandle::begin(
context,
prepared.context.clone(),
crate::trace::CaptureStart {
provider: prepared.harness.clone(),
model: prepared.model.clone(),
surface: "headless".to_string(),
input_op: "initial".to_string(),
gather_ms: prepared.context_gather_ms,
render_ms: prepared.context_render_ms,
raw_provider: true,
},
) {
Ok(capture) => Some(capture),
Err(err) => {
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
self.finish_failed_pass(
&body_id,
&format!("failed to establish trace capture: {err}"),
)
.await;
return;
}
},
None if cfg!(test) => None,
None => {
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
self.finish_failed_pass(&body_id, "trace capture identity is unavailable")
.await;
return;
}
};
let (event_tx, mut event_rx) = mpsc::unbounded_channel();
let (raw_tx, mut raw_rx) = mpsc::unbounded_channel();
let harness = match &self.backend {
BodyBackend::Harness { create, .. } => {
create(&prepared.harness, ApprovalPolicy::AutoApprove, event_tx)
}
#[cfg(test)]
BodyBackend::Process(_) => unreachable!("live skill requires a harness backend"),
};
let mut harness = match harness {
Ok(harness) => harness,
Err(err) => {
finish_capture(capture.as_ref(), "failed");
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
self.finish_failed_pass(
&body_id,
&format!("failed to create {} harness: {err:#}", prepared.harness),
)
.await;
return;
}
};
if capture.is_some() {
harness.set_raw_provider_sender(Some(raw_tx));
}
let resume_session_id =
provider_session_id_for_harness(self.provider_session.as_ref(), &prepared.harness);
if resume_session_id.is_none() {
self.provider_session = None;
}
harness.set_provider_session_id(resume_session_id);
if let Err(err) = harness.start(&prepared.config).await {
finish_capture(capture.as_ref(), "failed");
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
self.finish_failed_pass(
&body_id,
&format!("failed to start {} harness: {err:#}", prepared.harness),
)
.await;
return;
}
body.session_id = harness.provider_session_id();
if let Some(capture) = &capture {
capture.set_provider_session_id(body.session_id.clone());
}
if let Some(session_id) = &body.session_id {
self.provider_session = Some(ProviderSessionRef {
harness: prepared.harness.clone(),
session_id: session_id.clone(),
});
}
let mut body_session_id = body.session_id.clone();
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
if self.end.is_some() {
let _ = harness.stop().await;
finish_capture(capture.as_ref(), "interrupted");
return;
}
if let Err(err) = harness.send_input(&prepared.input).await {
let _ = harness.stop().await;
finish_capture(capture.as_ref(), "failed");
self.finish_failed_pass(
&body_id,
&format!(
"failed to start {} / {} turn: {err:#}",
step.flow, step.step
),
)
.await;
return;
}
let supports_steer = harness.capabilities().supports_steer;
let mut timeout = Box::pin(tokio::time::sleep(self.config.pass_timeout));
let mut terminal_wait = Box::pin(tokio::time::sleep(Duration::from_secs(86_400)));
let mut terminal_status: Option<Lifecycle> = None;
let mut usage = TurnUsage::default();
loop {
tokio::select! {
biased;
item = inbox_rx.recv() => {
match self.inbox_action(item) {
InboxAction::Interrupt { skip } => {
self.interrupt_harness(&body_id, harness.as_mut(), skip).await;
finish_capture(capture.as_ref(), "interrupted");
return;
}
InboxAction::Deliver(InboxItem::Message(message))
if message.op == MessageOp::Steer && supports_steer =>
{
if self
.steer_harness(message, harness.as_mut(), capture.as_ref())
.await
{
timeout.as_mut().reset(Instant::now() + self.config.pass_timeout);
}
}
InboxAction::Deliver(InboxItem::Message(message))
if message.op == MessageOp::Steer =>
{
if self.seen.insert(message.id.clone()) {
self.queue.push(message);
}
self.interrupt_harness(&body_id, harness.as_mut(), false).await;
return;
}
InboxAction::Deliver(item) => self.on_inbox(item).await,
InboxAction::ListenerGone => {
let _ = harness.stop().await;
finish_capture(capture.as_ref(), "interrupted");
return;
}
}
}
raw = raw_rx.recv(), if capture.is_some() => {
if let (Some(raw), Some(capture)) = (raw, capture.as_ref()) {
capture.record_raw(raw.stream, &raw.line);
}
}
event = event_rx.recv() => {
let Some(event) = event else {
let _ = harness.stop().await;
self.finish_failed_pass(&body_id, "harness event stream closed").await;
finish_capture(capture.as_ref(), "failed");
return;
};
if let Some(capture) = &capture {
capture.record_conversation(event.clone());
}
if body_session_id.is_none() {
tokio::task::yield_now().await;
if let Some(session_id) = harness.provider_session_id() {
body_session_id = Some(session_id.clone());
if let Some(capture) = &capture {
capture.set_provider_session_id(Some(session_id.clone()));
}
self.provider_session = Some(ProviderSessionRef {
harness: prepared.harness.clone(),
session_id: session_id.clone(),
});
self.send(vec![ResidentDelta::BodySessionUpdated {
body_id: body_id.clone(),
session_id,
}]).await;
}
}
match event {
ConversationEvent::TextDelta { content, .. } => {
self.send(vec![ResidentDelta::TurnText { text: content }]).await;
}
ConversationEvent::ItemCompleted { item, .. } => {
self.send(vec![ResidentDelta::TurnItem { item }]).await;
}
ConversationEvent::TurnUsage { usage: reported, .. } => {
usage = reported;
self.send(vec![ResidentDelta::TurnUsage {
input_tokens: Some(usage.input_tokens),
output_tokens: Some(usage.output_tokens),
cache_read_tokens: usage.cache_read_tokens,
}]).await;
if terminal_status.is_some() {
let status = terminal_status.take().expect("checked");
let outcome = if status == Lifecycle::Completed {
"completed"
} else if status == Lifecycle::Interrupted {
"interrupted"
} else {
"failed"
};
self.finish_harness_pass(
&body_id,
&step,
status,
usage.cost_usd,
harness.as_mut(),
).await;
finish_capture(capture.as_ref(), outcome);
return;
}
}
ConversationEvent::TurnCompleted { status, .. } => {
terminal_status = Some(status);
terminal_wait.as_mut().reset(
Instant::now() + Duration::from_millis(100)
);
}
ConversationEvent::Error { code, message } => {
let _ = harness.stop().await;
self.finish_failed_pass(
&body_id,
&format!("{code}: {message}"),
).await;
finish_capture(capture.as_ref(), "failed");
return;
}
ConversationEvent::TurnStarted { .. }
| ConversationEvent::ItemStarted { .. }
| ConversationEvent::ItemUpdated { .. }
| ConversationEvent::ReasoningDelta { .. }
| ConversationEvent::DiffUpdated { .. }
| ConversationEvent::SuggestedActions { .. }
| ConversationEvent::StatusChanged { .. } => {}
}
if self.end.is_some() {
let _ = harness.stop().await;
finish_capture(capture.as_ref(), "interrupted");
return;
}
}
_ = &mut terminal_wait, if terminal_status.is_some() => {
let status = terminal_status.take().expect("checked");
let outcome = if status == Lifecycle::Completed {
"completed"
} else if status == Lifecycle::Interrupted {
"interrupted"
} else {
"failed"
};
self.finish_harness_pass(
&body_id,
&step,
status,
usage.cost_usd,
harness.as_mut(),
).await;
finish_capture(capture.as_ref(), outcome);
return;
}
_ = &mut timeout => {
match self.timeout_action() {
TimeoutAction::End => {
let _ = harness.stop().await;
finish_capture(capture.as_ref(), "interrupted");
return;
}
TimeoutAction::Expire => {
let _ = harness.interrupt().await;
let _ = harness.stop().await;
self.finish_timed_out_pass(&body_id).await;
finish_capture(capture.as_ref(), "interrupted");
return;
}
}
}
}
}
}
async fn open_body(&mut self, body: BodyProvenance, answers: Vec<String>) {
self.send(vec![
ResidentDelta::BodyStarted { body },
ResidentDelta::TurnOpened { answers },
])
.await;
}
async fn steer_harness(
&mut self,
message: PendingMessage,
harness: &mut dyn Harness,
capture: Option<&crate::trace::CaptureHandle>,
) -> bool {
if !self.seen.insert(message.id.clone()) {
return false;
}
let id = message.id.0.clone();
self.send(vec![ResidentDelta::TurnSteered {
answers: vec![id.clone()],
}])
.await;
if self.end.is_some() {
return false;
}
let text = match &message.from {
Some(from) => format!("[{from}] {}", message.text),
None => message.text.clone(),
};
if let Some(capture) = capture {
if let Err(err) = capture.begin_turn("steer", &text) {
tracing::warn!(error = %err, "trace capture refused live steering; requeueing message");
self.send(vec![ResidentDelta::MessagesRequeued { ids: vec![id] }])
.await;
self.queue.push(message);
return false;
}
}
if let Err(err) = harness.send_input(&text).await {
tracing::warn!(error = %err, "live steering failed; requeueing message");
self.send(vec![ResidentDelta::MessagesRequeued { ids: vec![id] }])
.await;
self.queue.push(message);
return false;
}
true
}
async fn announce_interrupt(&mut self) {
self.send(vec![ResidentDelta::LoopState {
to: ResidentStateTo::Interrupting,
reason: "user interrupt".to_string(),
}])
.await;
}
async fn interrupt_harness(&mut self, body_id: &str, harness: &mut dyn Harness, skip: bool) {
self.announce_interrupt().await;
let _ = harness.interrupt().await;
let _ = harness.stop().await;
self.finish_interrupted_pass(body_id, skip).await;
}
fn inbox_action(&mut self, item: Option<InboxItem>) -> InboxAction {
match item {
Some(InboxItem::Interrupt) => InboxAction::Interrupt { skip: false },
Some(InboxItem::Skip) => InboxAction::Interrupt { skip: true },
Some(InboxItem::Message(message)) if message.op == MessageOp::Interrupt => {
if self.seen.insert(message.id.clone()) {
self.queue.push(message);
}
InboxAction::Interrupt { skip: false }
}
Some(InboxItem::Task(observation)) => InboxAction::Deliver(InboxItem::Message(
crate::wave::journal::task_observation_message(&observation),
)),
Some(item) => InboxAction::Deliver(item),
None => {
self.end = Some(LoopEnd::ListenerGone);
InboxAction::ListenerGone
}
}
}
fn timeout_action(&self) -> TimeoutAction {
if self.end.is_some() {
TimeoutAction::End
} else {
TimeoutAction::Expire
}
}
async fn finish_timed_out_pass(&mut self, body_id: &str) {
self.finish_failed_pass(
body_id,
&format!(
"wave timed out after {}s",
self.config.pass_timeout.as_secs()
),
)
.await;
}
async fn finish_harness_pass(
&mut self,
body_id: &str,
step: &StepRef,
status: Lifecycle,
cost_usd: Option<f64>,
harness: &mut dyn Harness,
) {
let _ = harness.stop().await;
match status {
Lifecycle::Completed => {
if let Err(err) =
crate::lf::commands::flow::commit_skill_work(&self.cwd, &step.step)
{
self.finish_failed_pass(
body_id,
&format!("failed to commit {}: {err:#}", step.step),
)
.await;
return;
}
self.consecutive_failures = 0;
self.finish_pass(body_id, StepOutcome::Completed, None, cost_usd)
.await;
}
Lifecycle::Interrupted => self.finish_interrupted_pass(body_id, false).await,
Lifecycle::Failed => {
self.finish_failed_pass(body_id, "harness turn failed")
.await
}
Lifecycle::Pending | Lifecycle::Running => {
self.finish_failed_pass(body_id, "harness ended without a terminal status")
.await;
}
}
}
async fn on_pass_output(
&mut self,
body_id: &str,
result: std::io::Result<std::process::Output>,
) {
match result {
Ok(output) if output.status.success() => {
self.consecutive_failures = 0;
self.ship_output(output).await;
self.finish_pass(body_id, StepOutcome::Completed, None, None)
.await;
}
Ok(output) => {
self.ship_output(output).await;
self.finish_failed_pass(body_id, "wave step exited nonzero")
.await;
}
Err(err) => {
self.finish_failed_pass(body_id, &format!("wave wait failed: {err:#}"))
.await;
}
}
}
async fn ship_output(&mut self, output: std::process::Output) {
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
let text = match (stdout.trim(), stderr.trim()) {
("", "") => String::new(),
(out, "") => out.to_string(),
("", err) => err.to_string(),
(out, err) => format!("{out}\n\nstderr:\n{err}"),
};
if !text.is_empty() {
self.send(vec![ResidentDelta::TurnText { text }]).await;
}
}
async fn interrupt_child(
&mut self,
body_id: &str,
wait_task: &mut tokio::task::JoinHandle<std::io::Result<std::process::Output>>,
skip: bool,
) {
self.announce_interrupt().await;
wait_task.abort();
self.finish_interrupted_pass(body_id, skip).await;
}
async fn finish_pass(
&mut self,
body_id: &str,
outcome: StepOutcome,
reason: Option<String>,
cost_usd: Option<f64>,
) {
let status = match outcome {
StepOutcome::Completed => Lifecycle::Completed,
StepOutcome::Skipped | StepOutcome::Interrupted => Lifecycle::Interrupted,
StepOutcome::Failed => Lifecycle::Failed,
};
let reason = reason.unwrap_or_else(|| outcome.name().to_string());
self.send(vec![
ResidentDelta::TurnFinished {
status,
cost_usd,
reason: (status != Lifecycle::Completed).then(|| reason.clone()),
},
ResidentDelta::BodyFinished {
body_id: body_id.to_string(),
outcome,
reason,
},
])
.await;
self.idle_since = Instant::now();
}
async fn finish_interrupted_pass(&mut self, body_id: &str, skip: bool) {
self.consecutive_failures = 0;
let (outcome, reason) = if skip {
(StepOutcome::Skipped, "skipped by user")
} else {
(StepOutcome::Interrupted, "interrupted by user")
};
self.finish_pass(body_id, outcome, Some(reason.to_string()), None)
.await;
}
async fn finish_failed_pass(&mut self, body_id: &str, reason: &str) {
self.finish_pass(body_id, StepOutcome::Failed, Some(reason.to_string()), None)
.await;
self.consecutive_failures += 1;
if self.consecutive_failures >= MAX_CONSECUTIVE_PASS_FAILURES {
self.fail(&format!(
"{MAX_CONSECUTIVE_PASS_FAILURES} consecutive wave failures: {reason}"
))
.await;
}
}
async fn fetch_context(&mut self) -> Option<crate::wave::wire::ContextResponse> {
match self.client.context().await {
Ok(context) => {
if self.provider_session.is_none() {
self.provider_session.clone_from(&context.provider_session);
}
Some(context)
}
Err(err) => {
tracing::info!(
error = %format!("{err:#}"),
"listener unreachable; ending residency"
);
self.end = Some(LoopEnd::ListenerGone);
None
}
}
}
async fn send(&mut self, deltas: Vec<ResidentDelta>) {
if self.end.is_some() || deltas.is_empty() {
return;
}
if let Err(err) = self.client.send_deltas(deltas).await {
tracing::info!(
error = %format!("{err:#}"),
"listener unreachable; ending residency"
);
self.end = Some(LoopEnd::ListenerGone);
}
}
async fn fail(&mut self, reason: &str) {
if self.end.is_some() {
return;
}
tracing::error!(
wave = self.wave,
reason,
"wave loop failed; reporting and exiting"
);
self.send(vec![ResidentDelta::LoopState {
to: ResidentStateTo::Failed,
reason: reason.to_string(),
}])
.await;
if self.end.is_none() {
self.end = Some(LoopEnd::Failed(reason.to_string()));
}
}
}
fn spawn_wave_step(
cwd: &Path,
step: &StepRef,
seed: &str,
max_turns: Option<u32>,
) -> std::io::Result<tokio::process::Child> {
let mut command = tokio::process::Command::from(lf_command());
command.arg("-b");
if let Some(max_turns) = max_turns {
command.arg("--max-turns").arg(max_turns.to_string());
}
command
.arg("__flow-step")
.arg(&step.flow)
.arg(step.index.to_string())
.arg(seed)
.current_dir(cwd)
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
command.spawn()
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Mutex};
use crate::chat::turns::{ChatRole, ChatTurn};
use crate::harness::Capabilities;
use crate::wave::journal::{journal_path, EventKind, Journal};
use crate::wave::runtime::WaveRuntime;
use crate::wave::server::{self, ResidentDoor};
use crate::wave::state::LoopState;
use async_trait::async_trait;
struct TestLoop {
runtime: Arc<WaveRuntime>,
seeds: Arc<Mutex<Vec<String>>>,
loop_task: tokio::task::JoinHandle<Result<()>>,
listener: Option<tokio::runtime::Runtime>,
_tmp: tempfile::TempDir,
}
impl Drop for TestLoop {
fn drop(&mut self) {
if let Some(rt) = self.listener.take() {
rt.shutdown_background();
}
}
}
impl TestLoop {
fn journal_events(&self) -> Vec<EventKind> {
let path = journal_path(self.runtime.repo_root(), "ship");
let (_, events) = Journal::open(&path).expect("read journal");
events.into_iter().map(|e| e.kind).collect()
}
fn pass_count(&self) -> usize {
self.seeds.lock().unwrap().len()
}
fn seed(&self, index: usize) -> String {
self.seeds.lock().unwrap()[index].clone()
}
}
fn test_config(heartbeat: Duration) -> LoopConfig {
LoopConfig {
heartbeat_idle: heartbeat,
pass_timeout: Duration::from_secs(5),
max_turns: None,
}
}
fn boot(
heartbeat: Duration,
script: &'static str,
) -> impl std::future::Future<Output = TestLoop> {
boot_in(tempfile::tempdir().expect("tempdir"), heartbeat, script)
}
async fn boot_in(
tmp: tempfile::TempDir,
heartbeat: Duration,
script: &'static str,
) -> TestLoop {
boot_with(tmp, test_config(heartbeat), script).await
}
async fn boot_with(
tmp: tempfile::TempDir,
config: LoopConfig,
script: &'static str,
) -> TestLoop {
let seeds = Arc::new(Mutex::new(Vec::new()));
let spawn_seeds = seeds.clone();
let spawn_pass: SpawnPass = Box::new(move |cwd, _step, seed, _max_turns| {
spawn_seeds.lock().unwrap().push(seed.to_string());
let mut command = tokio::process::Command::new("sh");
command
.arg("-c")
.arg(script)
.current_dir(cwd)
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
command.spawn()
});
boot_backend(tmp, config, BodyBackend::Process(spawn_pass), seeds).await
}
async fn boot_backend(
tmp: tempfile::TempDir,
config: LoopConfig,
backend: BodyBackend,
seeds: Arc<Mutex<Vec<String>>>,
) -> TestLoop {
let runtime =
WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open runtime");
let door = ResidentDoor::new("test-token");
let std_listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
std_listener.set_nonblocking(true).unwrap();
let addr = std_listener.local_addr().unwrap();
let app = server::router(
runtime.clone(),
door,
None,
None,
server::ShutdownDoor::new(),
);
let listener = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.enable_all()
.build()
.expect("listener runtime");
listener.spawn(async move {
let tcp = tokio::net::TcpListener::from_std(std_listener).expect("adopt listener");
axum::serve(tcp, app).await.ok();
});
let client = ListenerClient::new(addr.to_string(), "test-token".to_string());
let attach = client.attach(std::process::id()).await.expect("attach");
assert_eq!(attach.wave, "ship");
let (inbox_tx, inbox_rx) = mpsc::unbounded_channel();
tokio::spawn(crate::wave::resident::follow_inbox(
addr.to_string(),
inbox_tx,
));
let loop_task = tokio::spawn(run_loop_with(
client,
inbox_rx,
tmp.path().to_path_buf(),
tmp.path().to_path_buf(),
"ship".into(),
config,
backend,
));
TestLoop {
runtime,
seeds,
loop_task,
listener: Some(listener),
_tmp: tmp,
}
}
async fn wait_for(what: &str, cond: impl Fn() -> bool) {
for _ in 0..500 {
if cond() {
return;
}
tokio::time::sleep(Duration::from_millis(5)).await;
}
panic!("condition not met in time: {what}");
}
fn started_answers(events: &[EventKind]) -> Vec<Vec<MessageId>> {
events
.iter()
.filter_map(|kind| match kind {
EventKind::TurnStarted { answers, .. } => Some(answers.clone()),
_ => None,
})
.collect()
}
fn message_id(turn: &ChatTurn) -> MessageId {
let seq = turn.id.strip_prefix("turn-").expect("user turn id");
MessageId(format!("msg-{seq}"))
}
fn wake_of(seed: &str) -> String {
let start = seed.find("<wake>\n").expect("seed has a wake") + "<wake>\n".len();
let end = seed.find("\n</wake>").expect("seed closes the wake");
seed[start..end].to_string()
}
struct SteeringHarness {
events: mpsc::UnboundedSender<ConversationEvent>,
inputs: Arc<Mutex<Vec<String>>>,
supports_steer: bool,
}
#[async_trait]
impl Harness for SteeringHarness {
async fn start(&mut self, _config: &crate::engine::AgentConfig) -> Result<()> {
Ok(())
}
async fn send_input(&mut self, content: &str) -> Result<()> {
let mut inputs = self.inputs.lock().expect("inputs lock");
inputs.push(content.to_string());
let index = inputs.len();
drop(inputs);
if index == 1 {
let _ = self.events.send(ConversationEvent::TurnStarted {
turn_id: "vendor-turn".to_string(),
});
let _ = self.events.send(ConversationEvent::TextDelta {
turn_id: "vendor-turn".to_string(),
content: "hello".to_string(),
});
} else {
let _ = self.events.send(ConversationEvent::TextDelta {
turn_id: "vendor-turn".to_string(),
content: " world".to_string(),
});
let _ = self.events.send(ConversationEvent::TurnCompleted {
turn_id: "vendor-turn".to_string(),
status: Lifecycle::Completed,
});
let _ = self.events.send(ConversationEvent::TurnUsage {
turn_id: "vendor-turn".to_string(),
usage: TurnUsage {
input_tokens: 20,
output_tokens: 2,
..TurnUsage::default()
},
});
}
Ok(())
}
async fn interrupt(&mut self) -> Result<()> {
Ok(())
}
async fn stop(&mut self) -> Result<()> {
Ok(())
}
fn capabilities(&self) -> Capabilities {
Capabilities {
supports_steer: self.supports_steer,
}
}
fn provider_session_id(&self) -> Option<String> {
(!self.inputs.lock().expect("inputs lock").is_empty())
.then(|| "vendor-session".to_string())
}
}
#[tokio::test]
async fn steer_reaches_the_live_body_and_streams_into_one_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
let status = std::process::Command::new("git")
.args(["init", "--quiet"])
.current_dir(tmp.path())
.status()
.expect("git init");
assert!(status.success());
let inputs = Arc::new(Mutex::new(Vec::new()));
let harness_inputs = inputs.clone();
let backend = BodyBackend::Harness {
prepare: Box::new(|skill, seed, _wave, max_turns| {
Ok(crate::lf::commands::run::PreparedHarnessTurn {
config: crate::engine::AgentConfig {
agent: Some("fake".to_string()),
cwd: None,
max_turns,
..crate::engine::AgentConfig::default()
},
input: format!("{skill}\n{seed}"),
context: crate::trace::PreparedTurnContext::from_prompts(
"",
&format!("{skill}\n{seed}"),
),
harness: "fake".to_string(),
model: None,
context_gather_ms: 0,
context_render_ms: 0,
})
}),
create: Box::new(move |_name, _approval, events| {
Ok(Box::new(SteeringHarness {
events,
inputs: harness_inputs.clone(),
supports_steer: true,
}))
}),
};
let loop_ = boot_backend(
tmp,
test_config(Duration::from_secs(600)),
backend,
Arc::new(Mutex::new(Vec::new())),
)
.await;
let runtime = loop_.runtime.clone();
runtime
.deliver(MessageOp::Message, "begin".into())
.expect("user turn");
wait_for("initial live input", || inputs.lock().unwrap().len() == 1).await;
let steer = runtime
.deliver(MessageOp::Steer, "finish".into())
.expect("user turn");
wait_for("completed streamed turn", || {
runtime.thread_snapshot().iter().any(|turn| {
turn.role == ChatRole::Assistant
&& turn.status == Lifecycle::Completed
&& turn.text == "hello world"
})
})
.await;
assert_eq!(inputs.lock().unwrap()[1], "finish");
let completed = runtime
.thread_snapshot()
.into_iter()
.find(|turn| turn.role == ChatRole::Assistant)
.expect("assistant turn");
assert_eq!(
completed.body.and_then(|body| body.session_id),
Some("vendor-session".to_string())
);
assert!(loop_.journal_events().iter().any(|kind| matches!(
kind,
EventKind::TurnSteered { answers, .. }
if answers == &[message_id(&steer)]
)));
assert!(inputs.lock().unwrap()[0].contains("wave_clarify"));
}
#[tokio::test]
async fn unsupported_steer_restarts_the_same_wave_step() {
let tmp = tempfile::tempdir().expect("tempdir");
let status = std::process::Command::new("git")
.args(["init", "--quiet"])
.current_dir(tmp.path())
.status()
.expect("git init");
assert!(status.success());
let inputs = Arc::new(Mutex::new(Vec::new()));
let harness_inputs = inputs.clone();
let backend = BodyBackend::Harness {
prepare: Box::new(|skill, seed, _wave, max_turns| {
Ok(crate::lf::commands::run::PreparedHarnessTurn {
config: crate::engine::AgentConfig {
agent: Some("fake".to_string()),
cwd: None,
max_turns,
..crate::engine::AgentConfig::default()
},
input: format!("{skill}\n{seed}"),
context: crate::trace::PreparedTurnContext::from_prompts(
"",
&format!("{skill}\n{seed}"),
),
harness: "fake".to_string(),
model: None,
context_gather_ms: 0,
context_render_ms: 0,
})
}),
create: Box::new(move |_name, _approval, events| {
Ok(Box::new(SteeringHarness {
events,
inputs: harness_inputs.clone(),
supports_steer: false,
}))
}),
};
let loop_ = boot_backend(
tmp,
test_config(Duration::from_secs(600)),
backend,
Arc::new(Mutex::new(Vec::new())),
)
.await;
let runtime = loop_.runtime.clone();
runtime
.deliver(MessageOp::Message, "begin".into())
.expect("user turn");
wait_for("initial live input", || inputs.lock().unwrap().len() == 1).await;
runtime
.deliver(MessageOp::Steer, "finish differently".into())
.expect("steer");
wait_for("replacement turn", || inputs.lock().unwrap().len() >= 2).await;
let inputs = inputs.lock().unwrap();
assert!(inputs[0].starts_with("wave_clarify\n"));
assert!(inputs[1].starts_with("wave_clarify\n"));
assert!(inputs[1].contains("finish differently"));
assert!(runtime
.thread_snapshot()
.iter()
.any(|turn| turn.status == Lifecycle::Interrupted));
}
#[test]
fn path_for_children_starts_with_this_executables_dir() {
let exe_dir = std::env::current_exe()
.expect("current exe")
.parent()
.expect("exe has a dir")
.to_path_buf();
let path = path_for_children();
let first = std::env::split_paths(&path).next().expect("PATH non-empty");
assert_eq!(
first, exe_dir,
"the loop's PATH resolves `lf` to this build first"
);
}
#[test]
fn provider_sessions_resume_only_through_their_own_harness() {
let session = ProviderSessionRef {
harness: "codex".to_string(),
session_id: "thread-resume".to_string(),
};
assert_eq!(
provider_session_id_for_harness(Some(&session), "codex"),
Some("thread-resume".to_string())
);
assert_eq!(
provider_session_id_for_harness(Some(&session), "claude"),
None
);
}
#[test]
fn wave_pass_seed_carries_goal_loopflow_and_wake() {
let tmp = tempfile::tempdir().expect("tempdir");
let goal_dir = tmp.path().join("wave/ship");
std::fs::create_dir_all(&goal_dir).expect("goal dir");
std::fs::write(goal_dir.join("GOAL.md"), "Ship the thing.").expect("goal");
let seed = wave_pass_seed(tmp.path(), "ship", "hello from chat");
assert!(seed.contains("Ship the thing."));
assert!(seed.contains("<lf:loopflow>"));
assert!(seed.contains("<wake>\nhello from chat\n</wake>"));
assert_eq!(
LoopConfig::default().heartbeat_idle,
Duration::from_secs(4 * 60 * 60)
);
}
#[tokio::test]
async fn one_wake_runs_one_full_wave_flow_then_idles() {
let loop_ = boot(Duration::from_secs(600), "echo done").await;
loop_
.runtime
.deliver(MessageOp::Message, "first wake".into())
.expect("first wake");
wait_for("first Wave flow", || loop_.pass_count() == 3).await;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
loop_.pass_count(),
3,
"a completed Wave flow waits instead of starting another iteration"
);
loop_
.runtime
.deliver(MessageOp::Message, "second wake".into())
.expect("second wake");
wait_for("second Wave flow", || loop_.pass_count() == 6).await;
}
#[tokio::test]
async fn message_while_idle_starts_a_pass_answering_it() {
let loop_ = boot(Duration::from_secs(600), "echo hi!").await;
let user_turn = loop_
.runtime
.deliver(MessageOp::Message, "hello wave".into())
.expect("user turn");
wait_for("pass spawned", || loop_.pass_count() == 1).await;
assert_eq!(wake_of(&loop_.seed(0)), "hello wave");
wait_for("assistant turn", || {
loop_.runtime.thread_snapshot().iter().any(|t| {
t.role == ChatRole::Assistant
&& t.status == Lifecycle::Completed
&& t.text.contains("hi!")
})
})
.await;
let answers = started_answers(&loop_.journal_events());
assert_eq!(answers[0], vec![message_id(&user_turn)]);
assert!(answers.iter().skip(1).all(Vec::is_empty));
}
#[tokio::test]
async fn say_wakes_the_loop_and_is_consumed_by_the_next_pass() {
let loop_ = boot(Duration::from_secs(600), "echo noted").await;
let turn = loop_.runtime.deliver_say(
"implement run-1 finished: PR #7, one surprise".into(),
"worker".into(),
);
wait_for("pass spawned", || loop_.pass_count() == 1).await;
assert_eq!(
wake_of(&loop_.seed(0)),
"[worker] implement run-1 finished: PR #7, one surprise",
"the pass wake carries the byline"
);
wait_for("turn journaled", || {
!started_answers(&loop_.journal_events()).is_empty()
})
.await;
assert_eq!(
started_answers(&loop_.journal_events())[0],
vec![message_id(&turn)],
"the say emission is consumed like any queued message"
);
}
#[tokio::test]
async fn messages_during_a_pass_coalesce_into_one_boundary_pass() {
let loop_ = boot(Duration::from_secs(600), "sleep 0.4; echo done").await;
loop_
.runtime
.deliver(MessageOp::Message, "first".into())
.expect("user turn");
wait_for("pass 1 spawned", || loop_.pass_count() == 1).await;
let m2 = loop_
.runtime
.deliver(MessageOp::Message, "second".into())
.expect("user turn");
let m3 = loop_
.runtime
.deliver(MessageOp::Message, "third".into())
.expect("user turn");
wait_for("next iteration spawned", || loop_.pass_count() == 4).await;
let wake = wake_of(&loop_.seed(3));
assert!(wake.contains("second") && wake.contains("third"));
wait_for("next iteration TurnStarted journaled", || {
started_answers(&loop_.journal_events()).len() == 4
})
.await;
let answers = started_answers(&loop_.journal_events());
assert!(answers[1].is_empty());
assert!(answers[2].is_empty());
assert_eq!(answers[3], vec![message_id(&m2), message_id(&m3)]);
}
fn write_goal_with_crons(tmp: &std::path::Path, crons_yaml: &str) {
let dir = tmp.join("wave/ship");
std::fs::create_dir_all(&dir).expect("wave dir");
std::fs::write(
dir.join("GOAL.md"),
format!("---\ncrons:\n{crons_yaml}---\nShip.\n"),
)
.expect("write GOAL.md");
}
#[test]
fn next_cron_fire_honors_grace_last_fired_and_garbage() {
let now = Utc::now();
let due = next_cron_fire("0 0 * * * *", None, now).expect("hourly parses");
assert!(due <= now, "an occurrence inside the grace window is due");
let fired = next_cron_fire("0 0 * * * *", Some(now), now).expect("hourly parses");
assert!(fired > now, "a just-fired schedule waits for the next slot");
assert!(next_cron_fire("not-a-cron", None, now).is_none());
}
#[test]
fn cron_prompt_names_each_due_flow() {
let due = vec![
WaveCronDef {
flow: "qa".into(),
schedule: "* * * * * *".into(),
},
WaveCronDef {
flow: "wave-polish".into(),
schedule: "0 0 0 * * Mon *".into(),
},
];
assert_eq!(
cron_prompt(&due),
"cron due: qa — dispatch it\ncron due: wave-polish — dispatch it"
);
}
#[tokio::test]
async fn cron_due_opens_a_system_pass() {
let tmp = tempfile::tempdir().expect("tempdir");
write_goal_with_crons(tmp.path(), " - flow: qa\n schedule: '* * * * * *'\n");
let loop_ = boot_in(tmp, Duration::from_secs(600), "echo ok").await;
wait_for("cron pass spawned", || loop_.pass_count() >= 1).await;
assert_eq!(wake_of(&loop_.seed(0)), "cron due: qa — dispatch it");
}
#[tokio::test]
async fn cron_not_due_stays_quiet() {
let tmp = tempfile::tempdir().expect("tempdir");
write_goal_with_crons(
tmp.path(),
" - flow: qa\n schedule: '0 0 0 1 1 * 2099'\n",
);
let loop_ = boot_in(tmp, Duration::from_secs(600), "echo ok").await;
tokio::time::sleep(Duration::from_millis(400)).await;
assert_eq!(loop_.pass_count(), 0, "nothing due, nothing fired");
}
#[tokio::test]
async fn heartbeat_fires_when_idle_and_not_while_a_pass_runs() {
let loop_ = boot(Duration::from_millis(50), "sleep 0.3; echo beat").await;
wait_for("heartbeat pass", || loop_.pass_count() == 1).await;
assert_eq!(wake_of(&loop_.seed(0)), HEARTBEAT_PROMPT);
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(loop_.pass_count(), 1, "no heartbeat while a pass runs");
wait_for("next heartbeat", || loop_.pass_count() >= 2).await;
}
#[tokio::test]
async fn failure_cap_reports_failed_and_exits_the_resident() {
let mut loop_ = boot(Duration::from_millis(30), "exit 1").await;
wait_for("loop failed", || {
matches!(loop_.runtime.loop_state(), LoopState::Failed { .. })
})
.await;
let LoopState::Failed { reason } = loop_.runtime.loop_state() else {
unreachable!()
};
assert!(reason.contains("consecutive wave failures"), "{reason}");
assert!(
loop_.pass_count() >= MAX_CONSECUTIVE_PASS_FAILURES as usize,
"the cap took the full ladder"
);
let outcome = tokio::time::timeout(Duration::from_secs(5), &mut loop_.loop_task)
.await
.expect("loop task ends")
.expect("loop task not cancelled");
let err = outcome.expect_err("loop failure is an error exit");
assert!(err.to_string().contains("consecutive wave failures"));
}
#[tokio::test]
async fn pass_timeout_kills_the_child_and_fails_the_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
let config = LoopConfig {
heartbeat_idle: Duration::from_secs(600),
pass_timeout: Duration::from_millis(100),
max_turns: None,
};
let loop_ = boot_with(tmp, config, "sleep 30").await;
loop_
.runtime
.deliver(MessageOp::Message, "go".into())
.expect("user turn");
wait_for("pass spawned", || loop_.pass_count() == 1).await;
wait_for("turn failed", || {
loop_
.runtime
.thread_snapshot()
.iter()
.any(|t| t.role == ChatRole::Assistant && t.status == Lifecycle::Failed)
})
.await;
wait_for("back to idle", || {
loop_.runtime.loop_state() == LoopState::Idle
})
.await;
}
#[tokio::test]
async fn interrupt_kills_the_pass_and_finalizes_the_turn_interrupted() {
let loop_ = boot(Duration::from_secs(600), "sleep 30").await;
loop_
.runtime
.deliver(MessageOp::Message, "start".into())
.expect("user turn");
wait_for("pass spawned", || loop_.pass_count() == 1).await;
wait_for("turning", || loop_.runtime.loop_state().name() == "turning").await;
loop_.runtime.deliver_interrupt();
wait_for("idle again", || {
loop_.runtime.loop_state() == LoopState::Idle
})
.await;
let thread = loop_.runtime.thread_snapshot();
assert_eq!(
thread.last().unwrap().status,
Lifecycle::Interrupted,
"partial turn finalized as a value, not a crash"
);
let path: Vec<(String, String)> = loop_
.journal_events()
.iter()
.filter_map(|kind| match kind {
EventKind::LoopState { from, to, .. } => {
Some((from.name().to_string(), to.name().to_string()))
}
_ => None,
})
.collect();
assert_eq!(
path,
vec![
("idle".to_string(), "turning".to_string()),
("turning".to_string(), "interrupting".to_string()),
("interrupting".to_string(), "idle".to_string()),
]
);
}
#[tokio::test]
async fn listener_death_ends_the_resident_cleanly() {
let mut loop_ = boot(Duration::from_secs(600), "echo ok").await;
loop_
.listener
.take()
.expect("listener alive")
.shutdown_background();
let outcome = tokio::time::timeout(Duration::from_secs(30), &mut loop_.loop_task)
.await
.expect("loop task ends after listener death")
.expect("loop task not cancelled");
assert!(
outcome.is_ok(),
"listener death is a clean exit: {outcome:?}"
);
}
}