use std::collections::{BTreeSet, HashMap, HashSet};
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::str::FromStr;
use std::sync::Arc;
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};
use crate::durable::{RunLease, WorkRef};
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, SendCurrentOutcome};
use crate::store::{open_store, storage_config_from_env, Store};
use crate::wave::journal::{MessageDestination, 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");
}
}
async fn finish_resident_skill_work(resident_repo: &Path, wave: &str, skill: &str) -> Result<()> {
validate_resident_changes(resident_repo, wave)?;
if !crate::lf::commands::flow::commit_skill_work(resident_repo, skill)? {
return Ok(());
}
let branch = crate::engine::git::current_branch(resident_repo)?
.unwrap_or_else(|| "detached HEAD".to_string());
let head = crate::engine::git::rev_parse(resident_repo, "HEAD")?;
let resident_repo = resident_repo.to_path_buf();
let retained_path = resident_repo.display().to_string();
let title = format!("wave/{wave}: curate resident state");
let body = format!("Curates the inline GOAL.md and MEMORY.md state owned by Wave `{wave}`.");
let options = crate::ops::LandOptions {
strict: true,
local: false,
create_pr: true,
complete: false,
next_slug: None,
worktree: None,
commit_message: None,
pr_title: Some(title),
pr_body: Some(body),
agent: None,
};
tokio::task::spawn_blocking(move || {
crate::lf::commands::ops::land_repo(&resident_repo, &options, &crate::ops::NullProgress)
})
.await
.map_err(|error| anyhow!("resident delivery task failed: {error}"))?
.map_err(|error| {
anyhow!(
"resident delivery failed; retained {branch} at {head} in {}: {error:#}",
retained_path
)
})?;
Ok(())
}
fn validate_resident_changes(resident_repo: &Path, wave: &str) -> Result<()> {
let changed = resident_changed_paths(resident_repo)?;
let wave_root = Path::new("wave").join(wave);
let allowed = BTreeSet::from([wave_root.join("GOAL.md"), wave_root.join("MEMORY.md")]);
let unexpected = changed
.difference(&allowed)
.map(|path| path.display().to_string())
.collect::<Vec<_>>();
if unexpected.is_empty() {
return Ok(());
}
anyhow::bail!(
"Wave resident may only curate wave/{wave}/GOAL.md and wave/{wave}/MEMORY.md; unexpected changes: {}",
unexpected.join(", ")
)
}
fn resident_changed_paths(repo: &Path) -> Result<BTreeSet<PathBuf>> {
let mut changed = BTreeSet::new();
for args in [
vec!["diff", "--name-only", "-z", "HEAD", "--"],
vec!["ls-files", "--others", "--exclude-standard", "-z"],
] {
let output = Command::new("git")
.arg("-C")
.arg(repo)
.args(&args)
.output()?;
if !output.status.success() {
anyhow::bail!(
"git {} failed while validating resident changes: {}",
args.join(" "),
String::from_utf8_lossy(&output.stderr).trim()
);
}
changed.extend(
output
.stdout
.split(|byte| *byte == 0)
.filter(|path| !path.is_empty())
.map(|path| PathBuf::from(String::from_utf8_lossy(path).into_owned())),
);
}
Ok(changed)
}
fn read_crons(resident_repo: &Path, wave: &str) -> Vec<WaveCronDef> {
read_wave_config(resident_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(
resident_repo: &Path,
origin_repo: &Path,
wave: &str,
wake: &str,
metric_context: &str,
) -> String {
let seed = build_goal_seed(resident_repo, origin_repo, wave);
format!(
"{}\n\n{}\n\n{seed}\n\n{metric_context}\n\n<lf:wave-executive-loop>\n1. What is most important?\n2. What signals are arriving?\n3. What works?\n4. What does not?\n5. What is the current strategy?\n6. How should strategy adjust?\n\nTreat metrics as evidence, never as automatic KR completion or a composite Wave score.\n</lf:wave-executive-loop>\n\n<wake>\n{wake}\n</wake>",
crate::engine::prompt::loopflow_section(),
orchestration_discipline(wave),
)
}
async fn wave_metric_context(
control: Option<&WaveControl>,
origin_repo: &Path,
wave_name: &str,
) -> String {
let result = async {
let store = match control {
Some(control) => control.store.clone(),
None => Arc::new(
crate::store::open_existing_store()
.await
.ok_or_else(|| anyhow!("local registry is unavailable"))?,
),
};
let wave = match control {
Some(control) => {
let WorkRef::Wave(wave_id) = &control.lease.work else {
return Err(anyhow!("ambient Wave Run does not carry a Wave identity"));
};
store
.get_wave(wave_id)
.await?
.ok_or_else(|| anyhow!("ambient Wave is absent from the registry"))?
}
None => {
let locator = crate::wave::WaveLocator::discover(origin_repo, wave_name)?;
store
.get_wave_at(&locator)
.await?
.ok_or_else(|| anyhow!("wave/{wave_name} is absent from the registry"))?
}
};
crate::ops::metrics::stored_wave_metric_portfolio(
&store,
&wave,
time::OffsetDateTime::now_utc(),
)
.await
}
.await;
crate::ops::metrics::metric_prompt_section("metric-portfolio", result)
}
fn build_goal_seed(resident_repo: &Path, origin_repo: &Path, wave: &str) -> String {
let memory =
crate::engine::wave_context::gather_wave_memory_from(origin_repo, resident_repo, wave)
.unwrap_or_default();
match load_goal(wave, resident_repo) {
Ok(goal) => {
let ctx = GoalRenderContext {
flows: available_flow_names(origin_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 Tasks with `lf task status`, `follow-up`, \
`steer`, `interrupt`, `wait`, and `resume`. Each task owns one stable \
worktree; ordered PRs own its serial branches 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();
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(Box<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,
>;
use crate::harness::CreateHarness as CreateBodyHarness;
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>,
resident_repo: PathBuf,
origin_repo: PathBuf,
wave: String,
config: LoopConfig,
) -> Result<()> {
let control = wave_control(&wave).await?;
let prepare_origin = origin_repo.clone();
let prepare_resident = resident_repo.clone();
let backend = BodyBackend::Harness {
prepare: Box::new(move |skill, message, wave, max_turns| {
crate::lf::commands::run::prepare_wave_harness_turn(
skill,
message,
wave,
max_turns,
&prepare_origin,
&prepare_resident,
)
}),
create: Box::new(default_create_harness),
};
run_loop_with(
client,
inbox_rx,
resident_repo,
origin_repo,
wave,
config,
backend,
control,
)
.await
}
struct WaveControl {
store: Arc<Store>,
lease: RunLease,
}
async fn wave_control(wave: &str) -> Result<Option<WaveControl>> {
if std::env::var_os(crate::durable::RUN_LEASE_ENV).is_none()
&& std::env::var_os(crate::durable::RUN_CONTEXT_ENV).is_none()
{
return Ok(None);
}
let store = Arc::new(open_store(&storage_config_from_env()?).await?);
let lease = crate::ops::required_run_lease(&store).await?;
let WorkRef::Wave(wave_id) = &lease.work else {
return Err(anyhow!(
"ambient Run {} does not own Wave Work",
lease.run_id
));
};
let registered = store
.get_wave(wave_id)
.await?
.ok_or_else(|| anyhow!("Wave {wave_id} is absent from the control store"))?;
if registered.name() != wave {
return Err(anyhow!(
"ambient Run {} owns Wave '{}', not '{wave}'",
lease.run_id,
registered.name()
));
}
Ok(Some(WaveControl { store, lease }))
}
#[allow(clippy::too_many_arguments)]
async fn run_loop_with(
client: ListenerClient,
mut inbox_rx: mpsc::UnboundedReceiver<InboxItem>,
resident_repo: PathBuf,
origin_repo: PathBuf,
wave: String,
config: LoopConfig,
backend: BodyBackend,
control: Option<WaveControl>,
) -> Result<()> {
let ask_lane = control.as_ref().map(|control| {
crate::ops::ask::AskLane::new(control.lease.work.clone(), control.lease.clone())
});
let mut wave_loop = WaveLoop {
client,
resident_repo,
origin_repo,
wave,
config,
queue: Vec::new(),
evidence_queue: Vec::new(),
seen: HashSet::new(),
backend,
consecutive_failures: 0,
idle_since: Instant::now(),
cron_last_fired: HashMap::new(),
provider_session: None,
control,
ask_lane,
end: None,
};
let mut ask_poll = tokio::time::interval(Duration::from_millis(200));
ask_poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
while wave_loop.end.is_none() {
wave_loop.service_resolvers().await;
if !wave_loop.queue.is_empty() {
wave_loop.start_queued_pass(&mut inbox_rx).await;
continue;
}
if !wave_loop.evidence_queue.is_empty() {
wave_loop.start_evidence_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;
}
_ = ask_poll.tick(), if wave_loop.ask_lane.is_some() => {}
}
}
match wave_loop.end {
Some(LoopEnd::Failed(reason)) => Err(anyhow!(reason)),
_ => Ok(()),
}
}
struct WaveLoop {
client: ListenerClient,
resident_repo: PathBuf,
origin_repo: PathBuf,
wave: String,
config: LoopConfig,
backend: BodyBackend,
queue: Vec<PendingMessage>,
evidence_queue: Vec<PendingMessage>,
seen: HashSet<MessageId>,
consecutive_failures: u32,
idle_since: Instant,
cron_last_fired: HashMap<String, DateTime<Utc>>,
provider_session: Option<ProviderSessionRef>,
control: Option<WaveControl>,
ask_lane: Option<crate::ops::ask::AskLane>,
end: Option<LoopEnd>,
}
impl WaveLoop {
async fn service_resolvers(&mut self) {
let (Some(control), Some(ask_lane)) = (&self.control, self.ask_lane.as_mut()) else {
return;
};
if let Err(error) = ask_lane.reconcile(&control.store).await {
tracing::warn!(%error, "failed to reconcile Wave Ask lane");
}
}
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.resident_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.evidence_queue.push(message);
}
}
InboxItem::Project(observation) => {
let message = crate::wave::journal::project_observation_message(&observation);
if self.seen.insert(message.id.clone()) {
self.evidence_queue.push(message);
}
}
InboxItem::Promotion {
parent_wave_id,
parent,
} => {
let wake = crate::wave::PromotionWake {
parent_wave_id,
parent,
};
let message = crate::wave::journal::promotion_wake_message(&wake);
if self.seen.insert(message.id.clone()) {
self.evidence_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(),
MessageDestination::Local,
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.resident_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(), MessageDestination::Local, inbox_rx)
.await;
}
async fn start_queued_pass(&mut self, inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>) {
let messages = take_destination_prefix(&mut self.queue);
self.start_message_pass(messages, inbox_rx).await;
}
async fn start_evidence_pass(&mut self, inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>) {
let messages = std::mem::take(&mut self.evidence_queue);
self.start_message_pass(messages, inbox_rx).await;
}
async fn start_message_pass(
&mut self,
messages: Vec<PendingMessage>,
inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>,
) {
let answers: Vec<MessageId> = messages.iter().map(|m| m.id.clone()).collect();
let destination = messages
.first()
.map(PendingMessage::destination)
.unwrap_or(MessageDestination::Local);
let content = messages
.iter()
.map(|m| m.text.clone())
.collect::<Vec<_>>()
.join("\n\n");
self.run_pass(content, answers, destination, inbox_rx).await;
}
async fn capture_control(
&mut self,
provider: &str,
model: Option<&str>,
) -> Result<(
Option<crate::durable::Basis>,
Option<crate::trace::SupervisedInvocation>,
)> {
let Some(control) = &self.control else {
return Ok((None, None));
};
let epoch = control.store.current_epoch(&control.lease.work).await?;
let mut run = control
.store
.current_run(&control.lease.work)
.await?
.ok_or_else(|| anyhow!("Wave Run authority disappeared before Invocation"))?;
if run.id != control.lease.run_id {
anyhow::bail!(
"Wave Run {} was replaced before Invocation by {}",
control.lease.run_id,
run.id
);
}
let process_group = crate::engine::process::current_process_group_id()
.ok_or_else(|| anyhow!("Wave resident has no isolated process group"))?;
if run.state == crate::durable::RunState::Reserved {
let receipt = control
.store
.advance_run(
&control.lease,
crate::durable::RunAdvance::RunStarting {
containment: crate::durable::Containment::ProcessGroup {
id: i64::from(process_group),
},
cwd: self.resident_repo.clone(),
},
)
.await?;
let crate::durable::AdvanceReceipt::Run(started) = receipt else {
unreachable!("RunStarting returns a Run receipt")
};
run = started;
}
let receipt = control
.store
.advance_run(
&control.lease,
crate::durable::RunAdvance::InvocationStarting {
route: crate::durable::InvocationRoute {
provider: provider.to_string(),
model: model.map(str::to_string),
account_id: None,
},
surface: "headless".to_string(),
resume_token: None,
answer_ask_id: None,
},
)
.await?;
let crate::durable::AdvanceReceipt::Invocation(invocation) = receipt else {
unreachable!("InvocationStarting returns an Invocation receipt")
};
Ok((
Some(epoch.current_basis),
Some(crate::trace::SupervisedInvocation {
invocation_id: invocation.id,
supervising_run_id: run.id,
account_id: None,
resume_token: None,
}),
))
}
async fn run_pass(
&mut self,
wake: String,
answers: Vec<MessageId>,
destination: MessageDestination,
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 metric_context =
wave_metric_context(self.control.as_ref(), &self.origin_repo, &self.wave).await;
let seed = wave_pass_seed(
&self.resident_repo,
&self.origin_repo,
&self.wave,
&wake,
&metric_context,
);
let live_skill = step.kind == StepKind::Skill
&& matches!(&self.backend, BodyBackend::Harness { .. });
if live_skill {
self.run_harness_pass(step, seed, answers, &destination, 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.resident_repo);
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.resident_repo, &step, &seed, self.config.max_turns)
}
#[cfg(test)]
BodyBackend::Process(spawn) => {
spawn(&self.resident_repo, &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));
let mut ask_poll = tokio::time::interval(Duration::from_millis(200));
ask_poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
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;
}
}
}
_ = ask_poll.tick(), if self.ask_lane.is_some() => {
self.service_resolvers().await;
}
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>,
destination: &MessageDestination,
inbox_rx: &mut mpsc::UnboundedReceiver<InboxItem>,
) {
let mut body = body_provenance(&step, &self.resident_repo);
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 (basis, control) = match self
.capture_control(&prepared.harness, prepared.model.as_deref())
.await
{
Ok(control) => control,
Err(error) => {
let body_id = body.body_id.clone();
self.open_body(body, answers).await;
self.finish_failed_pass(
&body_id,
&format!("failed to establish Wave Run Invocation: {error}"),
)
.await;
return;
}
};
let capture = match crate::journal::trace_capture_context(
&self.resident_repo,
Some(step.flow.clone()),
Some(step.step.clone()),
) {
Ok(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,
basis,
supervision: control,
},
) {
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;
}
},
Err(_) if cfg!(test) => None,
Err(_) => {
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 mut timeout = Box::pin(tokio::time::sleep(self.config.pass_timeout));
let mut ask_poll = tokio::time::interval(Duration::from_millis(200));
ask_poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
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(item) => match *item {
InboxItem::Message(message) if message.op == MessageOp::Steer => {
if message.destination() == *destination {
if self
.steer_harness(message, harness.as_mut())
.await
{
timeout.as_mut().reset(Instant::now() + self.config.pass_timeout);
}
} else {
self.on_inbox(InboxItem::Message(message)).await;
}
}
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::UsageCheckpoint { .. } => {}
ConversationEvent::TurnCompleted { status, .. } => {
let outcome = if status == Lifecycle::Completed {
"completed"
} else if status == Lifecycle::Interrupted {
"interrupted"
} else {
"failed"
};
self.finish_harness_pass(
&body_id,
&step,
status,
harness.as_mut(),
).await;
finish_capture(capture.as_ref(), outcome);
return;
}
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 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;
}
}
}
_ = ask_poll.tick(), if self.ask_lane.is_some() => {
self.service_resolvers().await;
}
}
}
}
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) -> 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;
}
match harness.send_current(&message.text).await {
SendCurrentOutcome::Sent { .. } => {
if message.source.is_some() {
return true;
}
self.send(vec![ResidentDelta::MessagesRequeued { ids: vec![id] }])
.await;
self.queue.push(message);
true
}
SendCurrentOutcome::NotSteerable => {
self.send(vec![ResidentDelta::MessagesRequeued { ids: vec![id] }])
.await;
self.queue.push(message);
false
}
SendCurrentOutcome::Failed { error } | SendCurrentOutcome::Unknown { error, .. } => {
tracing::warn!(%error, "live steering was not confirmed; retaining message for next seed");
self.send(vec![ResidentDelta::MessagesRequeued { ids: vec![id] }])
.await;
self.queue.push(message);
false
}
}
}
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(item) => InboxAction::Deliver(Box::new(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,
harness: &mut dyn Harness,
) {
let _ = harness.stop().await;
match status {
Lifecycle::Completed => {
if let Err(err) =
finish_resident_skill_work(&self.resident_repo, &self.wave, &step.step).await
{
self.finish_failed_pass(
body_id,
&format!("failed to deliver {}: {err:#}", step.step),
)
.await;
return;
}
self.consecutive_failures = 0;
self.finish_pass(body_id, StepOutcome::Completed, None)
.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)
.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>) {
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,
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()))
.await;
}
async fn finish_failed_pass(&mut self, body_id: &str, reason: &str) {
self.finish_pass(body_id, StepOutcome::Failed, Some(reason.to_string()))
.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 take_destination_prefix(queue: &mut Vec<PendingMessage>) -> Vec<PendingMessage> {
let destination = queue
.first()
.expect("a queued pass starts with a message")
.destination();
let split = queue
.iter()
.position(|message| message.destination() != destination)
.unwrap_or(queue.len());
let remaining = queue.split_off(split);
std::mem::replace(queue, remaining)
}
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::chat::types::TurnUsage;
use crate::engine::OccurrencePolicy;
use crate::wave::journal::{
journal_path, DiscordChatBinding, DiscordMessageSource, EventKind, Journal,
};
use crate::wave::playhead::{Playhead, PlayheadEvent, QueuedInvocation, StepPlan};
use crate::wave::runtime::WaveRuntime;
use crate::wave::server::{self, ResidentDoor};
use crate::wave::state::LoopState;
use async_trait::async_trait;
fn queued_message(id: &str, source: Option<DiscordMessageSource>) -> PendingMessage {
PendingMessage {
id: MessageId(id.into()),
op: MessageOp::Message,
text: id.into(),
source,
}
}
fn discord_source(id: &str) -> DiscordMessageSource {
DiscordMessageSource {
binding: DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
},
message_id: id.into(),
author_id: "human".into(),
}
}
#[test]
fn resident_completion_accepts_only_its_wave_state_files() {
let repo = loopflow_test_support::TestRepo::new();
repo.create_file("wave/ship/GOAL.md", "Ship.\n");
repo.create_file("wave/ship/MEMORY.md", "Empty.\n");
repo.create_file("src/lib.rs", "// source\n");
repo.stage_all();
repo.commit("base");
repo.create_file("wave/ship/MEMORY.md", "Curated.\n");
validate_resident_changes(repo.path(), "ship").expect("owned memory is allowed");
repo.create_file("src/lib.rs", "// changed\n");
let error = validate_resident_changes(repo.path(), "ship")
.expect_err("source edits cannot auto-land");
assert!(error.to_string().contains("src/lib.rs"));
}
#[tokio::test]
async fn resident_completion_does_not_publish_an_unchanged_tree() {
let repo = loopflow_test_support::TestRepo::new();
repo.create_file("wave/ship/GOAL.md", "Ship.\n");
repo.stage_all();
repo.commit("base");
let head = crate::engine::git::rev_parse(repo.path(), "HEAD").unwrap();
finish_resident_skill_work(repo.path(), "ship", "assess")
.await
.unwrap();
assert_eq!(
crate::engine::git::rev_parse(repo.path(), "HEAD").unwrap(),
head
);
}
#[test]
fn discord_chat_batches_only_the_fifo_prefix_for_one_destination() {
let mut queue = vec![
queued_message("discord-1", Some(discord_source("1"))),
queued_message("discord-2", Some(discord_source("2"))),
queued_message("local", None),
queued_message("discord-3", Some(discord_source("3"))),
];
let first = take_destination_prefix(&mut queue);
assert_eq!(
first
.iter()
.map(|message| message.id.0.as_str())
.collect::<Vec<_>>(),
vec!["discord-1", "discord-2"]
);
assert_eq!(
queue
.iter()
.map(|message| message.id.0.as_str())
.collect::<Vec<_>>(),
vec!["local", "discord-3"]
);
let second = take_destination_prefix(&mut queue);
assert_eq!(second[0].id.0, "local");
assert_eq!(queue[0].id.0, "discord-3");
}
struct TestLoop {
runtime: Arc<WaveRuntime>,
seeds: Arc<Mutex<Vec<String>>>,
passes: mpsc::UnboundedReceiver<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()
}
async fn next_seed(&mut self) -> String {
self.passes.recv().await.expect("the loop spawns a pass")
}
}
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 (pass_tx, pass_rx) = mpsc::unbounded_channel();
let spawn_pass: SpawnPass = Box::new(move |cwd, _step, seed, _max_turns| {
spawn_seeds.lock().unwrap().push(seed.to_string());
let _ = pass_tx.send(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,
pass_rx,
None,
)
.await
}
async fn boot_backend(
tmp: tempfile::TempDir,
config: LoopConfig,
backend: BodyBackend,
seeds: Arc<Mutex<Vec<String>>>,
passes: mpsc::UnboundedReceiver<String>,
control: Option<WaveControl>,
) -> 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_with_observer(
runtime.clone(),
door,
Arc::new(crate::wave::registry::ObserverSlot::new(
runtime.clone(),
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,
control,
));
TestLoop {
runtime,
seeds,
passes,
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 init_test_git_repo(path: &Path) {
let status = std::process::Command::new("git")
.args([
"-c",
"user.name=Loopflow Test",
"-c",
"user.email=test@loopflow.dev",
"init",
"--quiet",
])
.current_dir(path)
.status()
.expect("git init");
assert!(status.success());
std::fs::write(path.join(".gitignore"), ".lf/\n").expect("write gitignore");
let status = std::process::Command::new("git")
.args(["add", ".gitignore"])
.current_dir(path)
.status()
.expect("git add");
assert!(status.success());
let status = std::process::Command::new("git")
.args([
"-c",
"user.name=Loopflow Test",
"-c",
"user.email=test@loopflow.dev",
"commit",
"--quiet",
"--allow-empty",
"-m",
"initial",
])
.current_dir(path)
.status()
.expect("git commit");
assert!(status.success());
}
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>>>,
accepts_current_send: bool,
}
struct CompletingHarness {
events: mpsc::UnboundedSender<ConversationEvent>,
inputs: Arc<Mutex<Vec<String>>>,
}
#[async_trait]
impl Harness for CompletingHarness {
async fn start(&mut self, _config: &crate::engine::AgentConfig) -> Result<()> {
Ok(())
}
async fn send_input(&mut self, content: &str) -> Result<()> {
self.inputs
.lock()
.expect("inputs lock")
.push(content.into());
let turn_id = "recovery-turn".to_string();
let _ = self.events.send(ConversationEvent::TurnStarted {
turn_id: turn_id.clone(),
});
let _ = self.events.send(ConversationEvent::TextDelta {
turn_id: turn_id.clone(),
content: "recovered".to_string(),
});
let _ = self.events.send(ConversationEvent::TurnCompleted {
turn_id,
status: Lifecycle::Completed,
});
Ok(())
}
async fn send_current(&mut self, _content: &str) -> SendCurrentOutcome {
SendCurrentOutcome::NotSteerable
}
async fn interrupt(&mut self) -> Result<()> {
Ok(())
}
async fn stop(&mut self) -> Result<()> {
Ok(())
}
fn provider_session_id(&self) -> Option<String> {
None
}
}
#[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::UsageCheckpoint {
turn_id: "vendor-turn".to_string(),
usage: TurnUsage {
input_tokens: Some(20),
output_tokens: Some(2),
..TurnUsage::default()
},
final_receipt: true,
});
let _ = self.events.send(ConversationEvent::TurnCompleted {
turn_id: "vendor-turn".to_string(),
status: Lifecycle::Completed,
});
}
Ok(())
}
async fn send_current(&mut self, content: &str) -> SendCurrentOutcome {
if !self.accepts_current_send {
return SendCurrentOutcome::NotSteerable;
}
match self.send_input(content).await {
Ok(()) => SendCurrentOutcome::Sent {
provider_turn_id: "vendor-turn".to_string(),
},
Err(error) => SendCurrentOutcome::Failed {
error: error.to_string(),
},
}
}
async fn interrupt(&mut self) -> Result<()> {
Ok(())
}
async fn stop(&mut self) -> Result<()> {
Ok(())
}
fn provider_session_id(&self) -> Option<String> {
(!self.inputs.lock().expect("inputs lock").is_empty())
.then(|| "vendor-session".to_string())
}
}
#[tokio::test]
async fn stale_wave_definition_resets_before_running_the_fresh_flow() {
let tmp = tempfile::tempdir().expect("tempdir");
init_test_git_repo(tmp.path());
let (mut journal, _) =
Journal::open(&journal_path(tmp.path(), "ship")).expect("open journal");
let root = QueuedInvocation {
id: "wave-root".to_string(),
flow: "wave".to_string(),
steps: ["wave_clarify", "wave_pursue", "wave_mutate"]
.into_iter()
.map(|name| StepPlan {
name: name.to_string(),
kind: StepKind::Skill,
policy: OccurrencePolicy::default(),
})
.collect(),
};
let (playhead, event) = Playhead::resume_root(root, 2, 7).expect("legacy playhead");
journal.append(|_| EventKind::PlayheadChanged {
event,
playhead: Box::new(playhead),
});
drop(journal);
let attempts = Arc::new(Mutex::new(Vec::new()));
let prepare_attempts = attempts.clone();
let inputs = Arc::new(Mutex::new(Vec::new()));
let harness_inputs = inputs.clone();
let backend = BodyBackend::Harness {
prepare: Box::new(move |skill, seed, _wave, max_turns| {
prepare_attempts
.lock()
.expect("attempts lock")
.push(skill.to_string());
Ok(crate::lf::commands::run::PreparedHarnessTurn {
config: crate::engine::AgentConfig {
agent: Some("fake".to_string()),
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(CompletingHarness {
events,
inputs: harness_inputs.clone(),
}))
}),
};
let loop_ = boot_backend(
tmp,
test_config(Duration::from_secs(600)),
backend,
Arc::new(Mutex::new(Vec::new())),
mpsc::unbounded_channel().1,
None,
)
.await;
let user_turn = loop_
.runtime
.deliver(MessageOp::Message, "recover".into())
.expect("recovery wake");
wait_for("fresh Wave iteration completed", || {
loop_.runtime.thread_snapshot().iter().any(|turn| {
turn.role == ChatRole::Assistant
&& turn.status == Lifecycle::Completed
&& turn.text == "recovered"
})
})
.await;
wait_for("next Wave iteration is idle", || {
loop_.runtime.loop_state() == LoopState::Idle
&& loop_
.runtime
.playhead()
.and_then(|playhead| playhead.now)
.is_some_and(|step| step.step == "wave/clarify" && step.iteration == 1)
})
.await;
assert_eq!(
*attempts.lock().expect("attempts lock"),
vec!["wave/clarify", "wave/pursue", "wave/mutate"]
);
assert_eq!(inputs.lock().expect("inputs lock").len(), 3);
let events = loop_.journal_events();
assert_eq!(
events
.iter()
.filter(|event| matches!(
event,
EventKind::PlayheadChanged {
event: PlayheadEvent::DefinitionReset,
..
}
))
.count(),
1
);
assert_eq!(
events
.iter()
.filter(|event| matches!(
event,
EventKind::PlayheadChanged {
event: PlayheadEvent::StepStarted { .. },
..
}
))
.count(),
3,
"reset itself never opens a body; the fresh three-step flow does"
);
assert!(!events.iter().any(|event| matches!(
event,
EventKind::PlayheadChanged {
event: PlayheadEvent::StepFinished {
outcome: StepOutcome::Failed,
..
},
..
}
)));
assert_eq!(
started_answers(&events),
vec![vec![message_id(&user_turn)], Vec::new(), Vec::new()]
);
assert!(!loop_
.runtime
.thread_snapshot()
.iter()
.any(|turn| turn.role == ChatRole::Assistant && turn.status == Lifecycle::Failed));
}
#[tokio::test]
async fn steer_reaches_the_live_body_and_streams_into_one_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
init_test_git_repo(tmp.path());
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(),
accepts_current_send: true,
}))
}),
};
let loop_ = boot_backend(
tmp,
test_config(Duration::from_secs(600)),
backend,
Arc::new(Mutex::new(Vec::new())),
mpsc::unbounded_channel().1,
None,
)
.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_waits_for_the_next_wave_boundary() {
let tmp = tempfile::tempdir().expect("tempdir");
init_test_git_repo(tmp.path());
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(),
accepts_current_send: false,
}))
}),
};
let loop_ = boot_backend(
tmp,
test_config(Duration::from_secs(600)),
backend,
Arc::new(Mutex::new(Vec::new())),
mpsc::unbounded_channel().1,
None,
)
.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("steer requeued", || {
loop_
.journal_events()
.iter()
.any(|event| matches!(event, EventKind::MessagesRequeued { .. }))
})
.await;
let inputs = inputs.lock().unwrap();
assert!(inputs[0].starts_with("wave/clarify\n"));
assert_eq!(
inputs.len(),
1,
"plain steering does not interrupt the Turn"
);
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(),
tmp.path(),
"ship",
"hello from chat",
"<lf:metric-portfolio>\n{\"metrics\":[],\"contract_issues\":[]}\n</lf:metric-portfolio>",
);
assert!(seed.contains("Ship the thing."));
assert!(seed.contains("<lf:loopflow>"));
assert!(seed.contains("<lf:metric-portfolio>"));
assert!(seed.contains("What is most important?"));
assert!(seed.contains("How should strategy adjust?"));
assert!(seed.contains("<wake>\nhello from chat\n</wake>"));
let doctrine_at = seed.find("<lf:loopflow>").expect("doctrine");
let goal_at = seed.find("Ship the thing.").expect("goal");
let wake_at = seed.find("<wake>").expect("wake");
assert!(doctrine_at < goal_at, "doctrine precedes the goal seed");
assert!(goal_at < wake_at, "wake closes the seed");
assert_eq!(
LoopConfig::default().heartbeat_idle,
Duration::from_secs(4 * 60 * 60)
);
}
#[test]
fn wave_pass_reads_mutable_state_from_resident_and_catalog_from_origin() {
let origin = tempfile::tempdir().unwrap();
let resident = tempfile::tempdir().unwrap();
std::fs::create_dir_all(origin.path().join("wave/ship")).unwrap();
std::fs::create_dir_all(origin.path().join(".lf/flows")).unwrap();
std::fs::create_dir_all(resident.path().join("wave/ship")).unwrap();
std::fs::write(
origin.path().join("wave/ship/MEMORY.md"),
"stale canonical memory\n",
)
.unwrap();
std::fs::write(
origin.path().join(".lf/flows/canonical.yaml"),
"- implement\n",
)
.unwrap();
std::fs::write(
resident.path().join("wave/ship/GOAL.md"),
"Use the current state.\n",
)
.unwrap();
std::fs::write(
resident.path().join("wave/ship/MEMORY.md"),
"curated resident memory\n",
)
.unwrap();
let seed = wave_pass_seed(resident.path(), origin.path(), "ship", "continue", "");
assert!(seed.contains("curated resident memory"));
assert!(!seed.contains("stale canonical memory"));
assert!(seed.contains("canonical"));
}
#[tokio::test]
async fn one_wake_runs_one_full_wave_flow_then_idles() {
let mut loop_ = boot(Duration::from_secs(600), "echo done").await;
loop_
.runtime
.deliver(MessageOp::Message, "first wake".into())
.expect("first wake");
for _ in 0..3 {
loop_.next_seed().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");
for _ in 0..3 {
loop_.next_seed().await;
}
}
#[tokio::test]
async fn typed_promotion_wake_runs_one_child_flow_without_user_speech() {
let mut loop_ = boot(Duration::from_secs(600), "echo promoted").await;
let wake = crate::wave::PromotionWake {
parent_wave_id: crate::id::WaveId::new(),
parent: "platform".to_string(),
};
assert!(loop_.runtime.deliver_promotion_wake(wake.clone()));
assert!(
!loop_.runtime.deliver_promotion_wake(wake.clone()),
"a replayed promotion signal is deduplicated before scheduling"
);
assert_eq!(wake_of(&loop_.next_seed().await), wake.prompt());
wait_for("one promoted Wave flow", || loop_.pass_count() == 3).await;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
loop_.pass_count(),
3,
"one promotion fact starts one three-step Wave flow"
);
let events = loop_.journal_events();
assert_eq!(
events
.iter()
.filter(|event| matches!(event, EventKind::PromotionObserved { .. }))
.count(),
1
);
assert!(
!events
.iter()
.any(|event| matches!(event, EventKind::UserMessage { .. })),
"machine promotion never enters the human thread door"
);
assert_eq!(
started_answers(&events)[0],
vec![MessageId(wake.inbox_id())]
);
}
#[tokio::test]
async fn message_while_idle_starts_a_pass_answering_it() {
let mut loop_ = boot(Duration::from_secs(600), "echo hi!").await;
let user_turn = loop_
.runtime
.deliver(MessageOp::Message, "hello wave".into())
.expect("user turn");
assert_eq!(wake_of(&loop_.next_seed().await), "hello wave");
loop_.next_seed().await;
assert!(
loop_.runtime.thread_snapshot().iter().any(|t| {
t.role == ChatRole::Assistant
&& t.status == Lifecycle::Completed
&& t.text.contains("hi!")
}),
"the pass answers into a completed assistant turn"
);
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 messages_during_a_pass_coalesce_into_one_boundary_pass() {
let mut loop_ = boot(Duration::from_secs(600), "sleep 0.4; echo done").await;
loop_
.runtime
.deliver(MessageOp::Message, "first".into())
.expect("user turn");
loop_.next_seed().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");
loop_.next_seed().await;
loop_.next_seed().await;
let wake = wake_of(&loop_.next_seed().await);
assert!(wake.contains("second") && wake.contains("third"));
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 non_catalog_harness_prepare_failure_counts_toward_cap() {
let attempts = Arc::new(Mutex::new(0_u32));
let prepare_attempts = attempts.clone();
let backend = BodyBackend::Harness {
prepare: Box::new(move |_skill, _seed, _wave, _max_turns| {
*prepare_attempts.lock().expect("attempts lock") += 1;
Err(anyhow!("provider configuration is invalid"))
}),
create: Box::new(|_name, _approval, _events| {
Err(anyhow!("prepare failure must not create a harness"))
}),
};
let mut loop_ = boot_backend(
tempfile::tempdir().expect("tempdir"),
test_config(Duration::from_millis(30)),
backend,
Arc::new(Mutex::new(Vec::new())),
mpsc::unbounded_channel().1,
None,
)
.await;
wait_for("prepare failures exhaust the cap", || {
matches!(loop_.runtime.loop_state(), LoopState::Failed { .. })
})
.await;
assert_eq!(
*attempts.lock().expect("attempts lock"),
MAX_CONSECUTIVE_PASS_FAILURES
);
let outcome = tokio::time::timeout(Duration::from_secs(5), &mut loop_.loop_task)
.await
.expect("loop task ends")
.expect("loop task not cancelled");
let error = outcome.expect_err("failure cap ends the resident");
assert!(error
.to_string()
.contains("provider configuration is invalid"));
}
#[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:?}"
);
}
}