use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use tokio::sync::broadcast;
use crate::chat::turns::{ChatRole, ChatTurn, TurnDelta};
use crate::chat::types::{ConversationItem, Lifecycle};
use crate::engine::wave_config::read_wave_config;
use crate::project_session::ProjectObservation;
use crate::receipt::Receipt;
use crate::security::sanitize_fs_component;
use crate::task::TaskObservation;
use crate::wave::channel::matches_prefix;
use crate::wave::journal::{
fold_thread, journal_path, project_observation_message, restore_pending,
task_observation_message, EventKind, Journal, MessageId, MessageOp, PendingMessage, Usage,
};
use crate::wave::memory::Memory;
use crate::wave::playhead::{
now_rfc3339, BodyProvenance, Playhead, PlayheadEvent, PlayheadView, QueuedInvocation,
StepOutcome,
};
use crate::wave::state::{can_transition, LoopState};
use crate::wave::wire::{ProviderSessionRef, ResidentDelta, ResidentStateTo};
const TURN_BROADCAST_CAPACITY: usize = 256;
const STATE_BROADCAST_CAPACITY: usize = 64;
const PLAYHEAD_BROADCAST_CAPACITY: usize = 64;
const MEMORY_BROADCAST_CAPACITY: usize = 64;
const INBOX_BROADCAST_CAPACITY: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ChannelRole {
Primary,
Child,
}
pub fn wave_channel_name(wave: &str) -> String {
sanitize_fs_component(wave)
}
pub fn channel_role(wave: &str, channel: &str) -> Option<ChannelRole> {
let family = wave_channel_name(wave);
if channel == wave || channel == family {
return Some(ChannelRole::Primary);
}
matches_prefix(channel, &family).then_some(ChannelRole::Child)
}
#[derive(Debug)]
pub struct TurnFrame {
pub turn: ChatTurn,
pub json: String,
}
impl TurnFrame {
fn share(turn: ChatTurn) -> Arc<Self> {
let json = serde_json::to_string(&turn).expect("ChatTurn serializes to JSON");
Arc::new(Self { turn, json })
}
}
#[derive(Debug)]
pub struct TurnDeltaFrame {
pub delta: TurnDelta,
pub json: String,
}
impl TurnDeltaFrame {
fn share(turn_id: String, item: ConversationItem) -> Arc<Self> {
let delta = TurnDelta { turn_id, item };
let json = serde_json::to_string(&delta).expect("TurnDelta serializes to JSON");
Arc::new(Self { delta, json })
}
}
#[derive(Debug, Clone)]
pub enum TurnBroadcast {
Whole(Arc<TurnFrame>),
Delta(Arc<TurnDeltaFrame>),
}
#[cfg(test)]
impl TurnBroadcast {
fn expect_whole(&self) -> &ChatTurn {
match self {
Self::Whole(frame) => &frame.turn,
Self::Delta(_) => panic!("expected a whole-turn frame, got a delta"),
}
}
fn expect_delta(&self) -> &TurnDelta {
match self {
Self::Delta(frame) => &frame.delta,
Self::Whole(_) => panic!("expected a delta frame, got a whole turn"),
}
}
}
#[derive(Debug, Clone)]
pub enum InboxItem {
Message(PendingMessage),
Task(TaskObservation),
Project(ProjectObservation),
Interrupt,
Skip,
}
#[derive(Debug)]
pub struct Subscription {
pub turns: Vec<ChatTurn>,
pub turn_rx: broadcast::Receiver<TurnBroadcast>,
pub state: LoopState,
pub state_rx: broadcast::Receiver<LoopState>,
pub playhead: Option<PlayheadView>,
pub playhead_rx: broadcast::Receiver<PlayheadView>,
pub memory_rx: broadcast::Receiver<String>,
pub memory_adds: Vec<String>,
pub memory_add_rx: broadcast::Receiver<String>,
pub pending: Vec<PendingMessage>,
pub tasks: HashMap<MessageId, TaskObservation>,
pub projects: HashMap<MessageId, ProjectObservation>,
pub inbox_rx: broadcast::Receiver<InboxItem>,
}
#[derive(Debug)]
struct OpenTurn {
turn: ChatTurn,
usage: Usage,
text_items: usize,
claims: Vec<MessageId>,
}
#[derive(Debug)]
struct Inner {
journal: Journal,
thread: Vec<ChatTurn>,
open: Option<OpenTurn>,
drop_deltas_until_opened: bool,
state: LoopState,
playhead: Option<Playhead>,
last_assistant_turn_id: Option<String>,
pending_messages: Vec<PendingMessage>,
messages: HashMap<MessageId, PendingMessage>,
tasks: HashMap<MessageId, TaskObservation>,
projects: HashMap<MessageId, ProjectObservation>,
memory_adds: Vec<String>,
}
#[derive(Debug)]
pub struct WaveRuntime {
name: String,
channel_name: String,
repo_root: PathBuf,
inner: Mutex<Inner>,
turn_tx: broadcast::Sender<TurnBroadcast>,
state_tx: broadcast::Sender<LoopState>,
playhead_tx: broadcast::Sender<PlayheadView>,
memory_tx: broadcast::Sender<String>,
memory_add_tx: broadcast::Sender<String>,
memory: Memory,
inbox_tx: broadcast::Sender<InboxItem>,
resident_expected: AtomicBool,
}
impl WaveRuntime {
pub fn open(name: String, repo_root: PathBuf) -> anyhow::Result<Arc<Self>> {
let (mut journal, events) = Journal::open(&journal_path(&repo_root, &name))?;
let mut fold = fold_thread(&events);
const ABANDONED: &str = "startup janitor: body abandoned by server restart";
let mut playhead = fold.playhead.take();
if let Some(state) = playhead.as_mut() {
if let Some(active) = state.active.clone() {
let events =
state.finish_body(&active.body_id, StepOutcome::Interrupted, ABANDONED)?;
for event in events {
journal.append(|_| EventKind::PlayheadChanged {
event,
playhead: Box::new(state.clone()),
});
}
}
}
for mut turn in fold.open {
let finished = journal.append(|_| EventKind::TurnFinished {
turn_id: turn.id.clone(),
status: Lifecycle::Failed,
usage: Usage::empty(),
termination_reason: Some(ABANDONED.to_string()),
});
turn.status = Lifecycle::Failed;
turn.close_body(finished.at_rfc3339(), Some(ABANDONED.to_string()));
fold.turns.push(turn);
}
let requeued = restore_pending(
&mut fold.pending_messages,
&fold.messages,
&fold.open_claims,
);
if !requeued.is_empty() {
journal.append(|_| EventKind::MessagesRequeued {
ids: requeued.clone(),
});
}
let last_assistant_turn_id = fold
.turns
.iter()
.rev()
.find(|turn| turn.role == ChatRole::Assistant)
.map(|turn| turn.id.clone());
let state = if fold.state == LoopState::Idle {
LoopState::Idle
} else {
journal.append(|_| EventKind::LoopState {
from: fold.state.clone(),
to: LoopState::Idle,
reason: "startup janitor: no live turn after restart".to_string(),
});
LoopState::Idle
};
let (turn_tx, _) = broadcast::channel(TURN_BROADCAST_CAPACITY);
let (state_tx, _) = broadcast::channel(STATE_BROADCAST_CAPACITY);
let (playhead_tx, _) = broadcast::channel(PLAYHEAD_BROADCAST_CAPACITY);
let (memory_tx, _) = broadcast::channel(MEMORY_BROADCAST_CAPACITY);
let (memory_add_tx, _) = broadcast::channel(MEMORY_BROADCAST_CAPACITY);
let (inbox_tx, _) = broadcast::channel(INBOX_BROADCAST_CAPACITY);
let memory = Memory::for_wave(&repo_root, &name);
Ok(Arc::new(Self {
channel_name: wave_channel_name(&name),
name,
repo_root,
inner: Mutex::new(Inner {
journal,
thread: fold.turns,
open: None,
drop_deltas_until_opened: false,
state,
playhead,
last_assistant_turn_id,
pending_messages: fold.pending_messages,
messages: fold.messages,
tasks: fold.tasks,
projects: fold.projects,
memory_adds: fold.memory_adds,
}),
turn_tx,
state_tx,
playhead_tx,
memory_tx,
memory_add_tx,
memory,
inbox_tx,
resident_expected: AtomicBool::new(false),
}))
}
pub fn name(&self) -> &str {
&self.name
}
pub fn channel_name(&self) -> &str {
&self.channel_name
}
pub fn repo_root(&self) -> &std::path::Path {
&self.repo_root
}
pub fn paused(&self) -> bool {
read_wave_config(&self.repo_root, &self.name)
.and_then(|config| config.paused)
.unwrap_or(false)
}
pub fn memory(&self) -> &Memory {
&self.memory
}
fn inner(&self) -> MutexGuard<'_, Inner> {
self.inner.lock().expect("wave runtime lock poisoned")
}
pub fn thread_snapshot(&self) -> Vec<ChatTurn> {
snapshot_tail_locked(&self.inner(), None)
}
pub fn thread_tail(&self, limit: Option<usize>) -> Vec<ChatTurn> {
snapshot_tail_locked(&self.inner(), limit)
}
pub fn thread_len(&self) -> usize {
let inner = self.inner();
inner.thread.len() + usize::from(inner.open.is_some())
}
pub fn latest_provider_session(&self) -> Option<ProviderSessionRef> {
let inner = self.inner();
if let Some(body) = inner
.open
.as_ref()
.and_then(|open| open.turn.body.as_ref())
.filter(|body| body_has_harness(body))
{
return provider_session_from_body(body);
}
inner
.thread
.iter()
.rev()
.filter_map(|turn| turn.body.as_ref())
.find(|body| body_has_harness(body))
.and_then(provider_session_from_body)
}
pub fn loop_state(&self) -> LoopState {
self.inner().state.clone()
}
pub fn playhead(&self) -> Option<PlayheadView> {
self.inner().playhead.as_ref().map(Playhead::view)
}
pub fn ensure_playhead(&self) -> anyhow::Result<PlayheadView> {
let mut inner = self.inner();
if let Some(playhead) = inner.playhead.as_ref() {
return Ok(playhead.view());
}
let root = QueuedInvocation::load(&self.repo_root, "wave")?;
let (playhead, event) = Playhead::new(root);
let view = playhead.view();
inner.journal.append(|_| EventKind::PlayheadChanged {
event,
playhead: Box::new(playhead.clone()),
});
inner.playhead = Some(playhead);
let _ = self.playhead_tx.send(view.clone());
Ok(view)
}
pub fn enqueue_flow(&self, flow: &str) -> anyhow::Result<PlayheadView> {
self.ensure_playhead()?;
let invocation = QueuedInvocation::load(&self.repo_root, flow)?;
let mut inner = self.inner();
let event = inner
.playhead
.as_mut()
.expect("ensure_playhead initialized it")
.enqueue(invocation)?;
self.journal_playhead_locked(&mut inner, vec![event])
}
pub fn start_body(&self, body: BodyProvenance) -> anyhow::Result<PlayheadView> {
self.ensure_playhead()?;
let mut inner = self.inner();
let event = inner
.playhead
.as_mut()
.expect("ensure_playhead initialized it")
.start_body(body)?;
self.journal_playhead_locked(&mut inner, vec![event])
}
fn finish_body(
&self,
body_id: &str,
outcome: StepOutcome,
reason: &str,
) -> anyhow::Result<PlayheadView> {
let mut inner = self.inner();
self.finish_body_locked(&mut inner, body_id, outcome, reason)
}
fn finish_body_locked(
&self,
inner: &mut Inner,
body_id: &str,
outcome: StepOutcome,
reason: &str,
) -> anyhow::Result<PlayheadView> {
let events = inner
.playhead
.as_mut()
.ok_or_else(|| anyhow::anyhow!("playhead is not initialized"))?
.finish_body(body_id, outcome, reason)?;
self.journal_playhead_locked(inner, events)
}
fn update_body_session(&self, body_id: &str, session_id: &str) -> anyhow::Result<PlayheadView> {
let mut inner = self.inner();
let event = inner
.playhead
.as_mut()
.ok_or_else(|| anyhow::anyhow!("playhead is not initialized"))?
.update_body_session(body_id, session_id)?;
let view = self.journal_playhead_locked(&mut inner, vec![event])?;
if let Some(open) = inner.open.as_mut() {
if let Some(body) = open
.turn
.body
.as_mut()
.filter(|body| body.body_id == body_id)
{
body.session_id = Some(session_id.to_string());
let _ = self
.turn_tx
.send(TurnBroadcast::Whole(TurnFrame::share(open.turn.clone())));
}
}
Ok(view)
}
pub fn skip_current(&self, reason: &str) -> anyhow::Result<PlayheadView> {
self.ensure_playhead()?;
let mut inner = self.inner();
let playhead = inner
.playhead
.as_mut()
.expect("ensure_playhead initialized it");
if playhead.active.is_some() {
return Err(anyhow::anyhow!("current step still has a live body"));
}
let step = playhead
.current()
.ok_or_else(|| anyhow::anyhow!("playhead has no current step"))?;
let mut body = BodyProvenance::for_step(&step, &self.repo_root);
let body_id = body.body_id.clone();
body.ended_at = Some(body.started_at.clone());
body.termination_reason = Some(reason.to_string());
let mut events = vec![playhead.start_body(body)?];
events.extend(playhead.finish_body(&body_id, StepOutcome::Skipped, reason)?);
self.journal_playhead_locked(&mut inner, events)
}
fn journal_playhead_locked(
&self,
inner: &mut Inner,
events: Vec<PlayheadEvent>,
) -> anyhow::Result<PlayheadView> {
let playhead = inner
.playhead
.clone()
.ok_or_else(|| anyhow::anyhow!("playhead is not initialized"))?;
for event in events {
inner.journal.append(|_| EventKind::PlayheadChanged {
event,
playhead: Box::new(playhead.clone()),
});
}
let view = playhead.view();
let _ = self.playhead_tx.send(view.clone());
Ok(view)
}
pub fn pending_messages(&self) -> Vec<PendingMessage> {
self.inner().pending_messages.clone()
}
pub fn subscribe_inbox(&self) -> broadcast::Receiver<InboxItem> {
self.inbox_tx.subscribe()
}
pub fn subscribe_states(&self) -> broadcast::Receiver<LoopState> {
self.state_tx.subscribe()
}
pub fn subscribe_turns(&self) -> broadcast::Receiver<TurnBroadcast> {
self.turn_tx.subscribe()
}
pub fn resident_expected(&self) -> bool {
self.resident_expected.load(Ordering::Relaxed)
}
pub fn set_resident_expected(&self) {
self.resident_expected.store(true, Ordering::Relaxed);
}
pub fn in_family(&self, channel: &str) -> bool {
channel_role(&self.name, channel).is_some()
}
pub fn is_primary(&self, channel: &str) -> bool {
channel_role(&self.name, channel) == Some(ChannelRole::Primary)
}
pub fn update_memory(&self, content: &str, summary: &str) -> std::io::Result<()> {
let mut inner = self.inner();
self.memory.write(content)?;
inner.journal.append(|_| EventKind::MemoryUpdated {
summary: summary.to_string(),
});
inner.memory_adds.clear();
let _ = self.memory_tx.send(summary.to_string());
Ok(())
}
pub fn append_memory(&self, fact: &str, receipts: Vec<Receipt>) -> std::io::Result<()> {
let mut inner = self.inner();
inner.journal.append(|_| EventKind::MemoryAdded {
fact: fact.to_string(),
receipts,
});
inner.memory_adds.push(fact.to_string());
let _ = self.memory_add_tx.send(fact.to_string());
Ok(())
}
pub fn memory_adds(&self) -> Vec<String> {
self.inner().memory_adds.clone()
}
pub fn journal_server_started(&self, pid: u32, endpoint: &str) {
let mut inner = self.inner();
inner.journal.append(|_| EventKind::ServerStarted {
pid,
endpoint: endpoint.to_string(),
});
}
pub fn subscribe_with_snapshot(&self, limit: Option<usize>) -> Subscription {
let inner = self.inner();
Subscription {
turns: snapshot_tail_locked(&inner, limit),
turn_rx: self.turn_tx.subscribe(),
state: inner.state.clone(),
state_rx: self.state_tx.subscribe(),
playhead: inner.playhead.as_ref().map(Playhead::view),
playhead_rx: self.playhead_tx.subscribe(),
memory_rx: self.memory_tx.subscribe(),
memory_adds: inner.memory_adds.clone(),
memory_add_rx: self.memory_add_tx.subscribe(),
pending: inner.pending_messages.clone(),
tasks: inner.tasks.clone(),
projects: inner.projects.clone(),
inbox_rx: self.inbox_tx.subscribe(),
}
}
pub fn transition(&self, to: LoopState, reason: &str) -> bool {
let mut inner = self.inner();
self.transition_locked(&mut inner, to, reason)
}
fn transition_locked(&self, inner: &mut Inner, to: LoopState, reason: &str) -> bool {
if !can_transition(&inner.state, &to) {
tracing::warn!(
from = inner.state.name(),
to = to.name(),
reason,
"illegal loop-state transition refused"
);
return false;
}
let from = std::mem::replace(&mut inner.state, to.clone());
inner.journal.append(|_| EventKind::LoopState {
from,
to: to.clone(),
reason: reason.to_string(),
});
let _ = self.state_tx.send(to);
true
}
fn begin_interrupt(&self, reason: &str) -> bool {
let mut inner = self.inner();
let LoopState::Turning { turn_id } = inner.state.clone() else {
return false;
};
self.transition_locked(&mut inner, LoopState::Interrupting { turn_id }, reason)
}
pub fn force_finalize_open_turn(&self, status: Lifecycle, reason: &str) -> bool {
let mut inner = self.inner();
let open = inner.open.take();
let active_body_id = inner
.playhead
.as_ref()
.and_then(|playhead| playhead.active.as_ref())
.map(|body| body.body_id.clone());
if open.is_none() && active_body_id.is_none() {
return false;
}
if let Some(OpenTurn {
mut turn, claims, ..
}) = open
{
inner.drop_deltas_until_opened = true;
let finished = inner.journal.append(|_| EventKind::TurnFinished {
turn_id: turn.id.clone(),
status,
usage: Usage::empty(),
termination_reason: Some(reason.to_string()),
});
if status != Lifecycle::Completed {
self.requeue_locked(&mut inner, &claims);
}
turn.status = status;
turn.close_body(finished.at_rfc3339(), Some(reason.to_string()));
self.transition_locked(&mut inner, LoopState::Idle, reason);
self.commit_locked(&mut inner, turn);
}
if let Some(body_id) = active_body_id {
let outcome = match status {
Lifecycle::Interrupted => StepOutcome::Interrupted,
Lifecycle::Completed => StepOutcome::Completed,
Lifecycle::Pending | Lifecycle::Running | Lifecycle::Failed => StepOutcome::Failed,
};
self.finish_body_locked(&mut inner, &body_id, outcome, reason)
.expect("the active body belongs to an initialized playhead");
}
true
}
fn requeue_locked(&self, inner: &mut Inner, ids: &[MessageId]) {
let restored = restore_pending(&mut inner.pending_messages, &inner.messages, ids);
if restored.is_empty() {
return;
}
inner
.journal
.append(|_| EventKind::MessagesRequeued { ids: restored });
}
fn commit_locked(&self, inner: &mut Inner, turn: ChatTurn) -> ChatTurn {
turn.validate()
.expect("Wave thread entries must satisfy the ChatTurn wire invariant");
if turn.role == ChatRole::Assistant {
inner.last_assistant_turn_id = Some(turn.id.clone());
}
inner.thread.push(turn.clone());
let _ = self
.turn_tx
.send(TurnBroadcast::Whole(TurnFrame::share(turn.clone())));
turn
}
pub fn deliver(&self, op: MessageOp, text: String) -> Option<ChatTurn> {
if op == MessageOp::Interrupt && text.trim().is_empty() {
self.deliver_interrupt();
return None;
}
Some(self.deliver_message(text, op, None))
}
pub fn deliver_say(&self, text: String, from: String) -> ChatTurn {
self.deliver_message(text, MessageOp::Say, Some(from))
}
fn deliver_message(&self, text: String, op: MessageOp, from: Option<String>) -> ChatTurn {
let mut inner = self.inner();
let event = inner.journal.append(|seq| EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op,
text: text.clone(),
from: from.clone(),
});
let id = MessageId(format!("msg-{}", event.seq));
let mut turn = ChatTurn::user(format!("turn-{}", event.seq), text.clone());
turn.created_at = event.at_rfc3339();
turn.from = from.clone();
let turn = self.commit_locked(&mut inner, turn);
let pending = PendingMessage { id, op, text, from };
inner.messages.insert(pending.id.clone(), pending.clone());
inner.pending_messages.push(pending.clone());
let _ = self.inbox_tx.send(InboxItem::Message(pending));
turn
}
pub fn deliver_interrupt(&self) {
let _ = self.inbox_tx.send(InboxItem::Interrupt);
}
pub fn deliver_skip(&self) {
let _ = self.inbox_tx.send(InboxItem::Skip);
}
pub fn deliver_task_observation(&self, observation: TaskObservation) -> bool {
let mut inner = self.inner();
let pending = task_observation_message(&observation);
if inner.tasks.contains_key(&pending.id) {
return false;
}
let event = inner.journal.append(|_| EventKind::TaskObserved {
observation: observation.clone(),
});
let turn = ChatTurn::child_activity(
format!("turn-{}", event.seq),
event.at_rfc3339(),
"task".to_string(),
crate::chat::turns::ChildControlActivity::from_task(&observation),
);
self.commit_locked(&mut inner, turn);
inner.messages.insert(pending.id.clone(), pending.clone());
inner.tasks.insert(pending.id.clone(), observation.clone());
inner.pending_messages.push(pending);
let _ = self.inbox_tx.send(InboxItem::Task(observation));
true
}
pub fn deliver_project_observation(&self, observation: ProjectObservation) -> bool {
let mut inner = self.inner();
let pending = project_observation_message(&observation);
if inner.projects.contains_key(&pending.id) {
return false;
}
let event = inner.journal.append(|_| EventKind::ProjectObserved {
observation: observation.clone(),
});
let turn = ChatTurn::child_activity(
format!("turn-{}", event.seq),
event.at_rfc3339(),
"project".to_string(),
crate::chat::turns::ChildControlActivity::from_project(&observation),
);
self.commit_locked(&mut inner, turn);
inner.messages.insert(pending.id.clone(), pending.clone());
inner
.projects
.insert(pending.id.clone(), observation.clone());
inner.pending_messages.push(pending);
let _ = self.inbox_tx.send(InboxItem::Project(observation));
true
}
pub fn append_finalized_turn(&self, turn: ChatTurn, answers: Vec<MessageId>) -> ChatTurn {
let mut inner = self.inner();
let started = inner.journal.append(|seq| EventKind::TurnStarted {
turn_id: format!("turn-{seq}"),
answers,
body: turn.body.clone().map(Box::new),
});
let turn_id = format!("turn-{}", started.seq);
if !turn.text.is_empty() {
inner.journal.append(|_| EventKind::TurnItem {
turn_id: turn_id.clone(),
item: ConversationItem::Message {
id: "text-0".to_string(),
text: turn.text.clone(),
phase: None,
},
});
}
for item in &turn.items {
inner.journal.append(|_| EventKind::TurnItem {
turn_id: turn_id.clone(),
item: item.clone(),
});
}
inner.journal.append(|_| EventKind::TurnFinished {
turn_id: turn_id.clone(),
status: turn.status,
usage: Usage::empty(),
termination_reason: None,
});
let committed = ChatTurn {
id: turn_id,
created_at: started.at_rfc3339(),
..turn
};
self.commit_locked(&mut inner, committed)
}
pub fn apply_resident_delta(&self, delta: ResidentDelta) {
match delta {
ResidentDelta::TurnOpened { answers } => self.resident_turn_opened(answers),
ResidentDelta::TurnText { text } => self.resident_turn_text(text),
ResidentDelta::TurnItem { item } => self.resident_turn_item(item),
ResidentDelta::TurnUsage {
input_tokens,
output_tokens,
cache_read_tokens,
} => self.resident_turn_usage(input_tokens, output_tokens, cache_read_tokens),
ResidentDelta::TurnFinished {
status,
cost_usd,
reason,
} => self.resident_turn_finished(status, cost_usd, reason),
ResidentDelta::TurnSteered { answers } => self.resident_turn_steered(answers),
ResidentDelta::MessagesRequeued { ids } => self.resident_requeue(ids),
ResidentDelta::BodyStarted { body } => {
if let Err(err) = self.start_body(body) {
tracing::warn!(error = %err, "resident body start rejected");
}
}
ResidentDelta::BodySessionUpdated {
body_id,
session_id,
} => {
if let Err(err) = self.update_body_session(&body_id, &session_id) {
tracing::warn!(error = %err, "resident body session update rejected");
}
}
ResidentDelta::BodyFinished {
body_id,
outcome,
reason,
} => {
if let Err(err) = self.finish_body(&body_id, outcome, &reason) {
tracing::warn!(error = %err, "resident body finish rejected");
}
}
ResidentDelta::LoopState { to, reason } => match to {
ResidentStateTo::Interrupting => {
if !self.begin_interrupt(&reason) {
tracing::warn!(reason, "resident reported Interrupting with no live turn");
}
}
ResidentStateTo::Failed => {
self.transition(
LoopState::Failed {
reason: reason.clone(),
},
&reason,
);
}
},
}
}
fn resident_turn_opened(&self, answers: Vec<String>) {
let paused = self.paused();
let mut inner = self.inner();
inner.drop_deltas_until_opened = false;
if let Some(OpenTurn {
turn: mut stale,
usage,
claims,
..
}) = inner.open.take()
{
tracing::warn!(
turn_id = stale.id,
"TurnOpened over an open turn; closing the stale turn as failed"
);
inner.journal.append(|_| EventKind::TurnFinished {
turn_id: stale.id.clone(),
status: Lifecycle::Failed,
usage,
termination_reason: Some("stale open turn closed".to_string()),
});
self.requeue_locked(&mut inner, &claims);
stale.status = Lifecycle::Failed;
self.transition_locked(&mut inner, LoopState::Idle, "stale open turn closed");
self.commit_locked(&mut inner, stale);
}
if paused {
tracing::warn!(
wave = self.name,
"wave is paused (GOAL.md frontmatter); turn refused, deltas dropped until the next TurnOpened"
);
inner.drop_deltas_until_opened = true;
return;
}
let answers = claim_answers(&mut inner, answers);
let claims = answers.clone();
let body = inner
.playhead
.as_ref()
.and_then(|playhead| playhead.active.clone());
let event = inner.journal.append(|seq| EventKind::TurnStarted {
turn_id: format!("turn-{seq}"),
answers,
body: body.clone().map(Box::new),
});
let turn_id = format!("turn-{}", event.seq);
self.transition_locked(
&mut inner,
LoopState::Turning {
turn_id: turn_id.clone(),
},
"turn opened",
);
let open = ChatTurn {
id: turn_id,
role: ChatRole::Assistant,
text: String::new(),
status: Lifecycle::Running,
items: Vec::new(),
created_at: event.at_rfc3339(),
from: None,
body,
activity: None,
};
let _ = self
.turn_tx
.send(TurnBroadcast::Whole(TurnFrame::share(open.clone())));
inner.open = Some(OpenTurn {
turn: open,
usage: Usage::empty(),
text_items: 0,
claims,
});
}
fn resident_turn_text(&self, text: String) {
let mut inner = self.inner();
if inner.drop_deltas_until_opened {
return;
}
let Some(open) = inner.open.as_mut() else {
tracing::warn!("text delta with no open turn; dropped");
return;
};
let item = ConversationItem::Message {
id: format!("text-{}", open.text_items),
text,
phase: Some("stream".to_string()),
};
open.text_items += 1;
self.append_turn_item_locked(&mut inner, item);
}
fn resident_turn_item(&self, item: ConversationItem) {
if matches!(&item, ConversationItem::Thought { text, .. } if text.trim().is_empty()) {
return;
}
let mut inner = self.inner();
if inner.drop_deltas_until_opened {
return;
}
if inner.open.is_none() {
tracing::warn!("item delta with no open turn; dropped");
return;
}
self.append_turn_item_locked(&mut inner, item);
}
fn append_turn_item_locked(&self, inner: &mut Inner, item: ConversationItem) {
let open = &mut inner.open.as_mut().expect("checked by callers").turn;
let turn_id = open.id.clone();
open.absorb_item(item.clone());
let frame = TurnDeltaFrame::share(turn_id.clone(), item.clone());
inner
.journal
.append(|_| EventKind::TurnItem { turn_id, item });
let _ = self.turn_tx.send(TurnBroadcast::Delta(frame));
}
fn resident_turn_usage(
&self,
input_tokens: Option<u64>,
output_tokens: Option<u64>,
cache_read_tokens: Option<u64>,
) {
let mut inner = self.inner();
if inner.drop_deltas_until_opened {
return;
}
let Some(open) = inner.open.as_mut() else {
return;
};
open.usage.input_tokens = add_opt(open.usage.input_tokens, input_tokens);
open.usage.output_tokens = add_opt(open.usage.output_tokens, output_tokens);
open.usage.cache_read_tokens = add_opt(open.usage.cache_read_tokens, cache_read_tokens);
}
fn resident_turn_finished(
&self,
status: Lifecycle,
cost_usd: Option<f64>,
reason: Option<String>,
) {
let mut inner = self.inner();
if inner.drop_deltas_until_opened {
tracing::debug!("late TurnFinished after a force-finalize; dropped");
return;
}
let Some(OpenTurn {
mut turn,
mut usage,
claims,
..
}) = inner.open.take()
else {
tracing::warn!("TurnFinished with no open turn; dropped");
return;
};
usage.cost_usd = cost_usd;
inner.journal.append(|_| EventKind::TurnFinished {
turn_id: turn.id.clone(),
status,
usage,
termination_reason: reason.clone(),
});
if status != Lifecycle::Completed {
self.requeue_locked(&mut inner, &claims);
}
turn.status = status;
turn.close_body(now_rfc3339(), reason);
self.transition_locked(&mut inner, LoopState::Idle, "turn finalized");
self.commit_locked(&mut inner, turn);
}
fn resident_turn_steered(&self, answers: Vec<String>) {
let mut inner = self.inner();
let (turn_id, turn_live) = match inner.state.clone() {
LoopState::Turning { turn_id } | LoopState::Interrupting { turn_id } => (turn_id, true),
_ => match inner.last_assistant_turn_id.clone() {
Some(turn_id) => (turn_id, false),
None => {
tracing::warn!("TurnSteered with no assistant turn anywhere; kept pending");
return;
}
},
};
let answers = claim_answers(&mut inner, answers);
if answers.is_empty() {
return;
}
if turn_live {
if let Some(open) = inner.open.as_mut() {
open.claims.extend(answers.iter().cloned());
}
}
inner
.journal
.append(|_| EventKind::TurnSteered { turn_id, answers });
}
fn resident_requeue(&self, ids: Vec<String>) {
let mut inner = self.inner();
let ids: Vec<MessageId> = ids.into_iter().map(MessageId).collect();
if let Some(open) = inner.open.as_mut() {
open.claims.retain(|claim| !ids.contains(claim));
}
self.requeue_locked(&mut inner, &ids);
}
}
fn provider_session_from_body(body: &BodyProvenance) -> Option<ProviderSessionRef> {
let harness = body.harness.as_deref()?.trim();
let session_id = body.session_id.as_deref()?.trim();
if harness.is_empty() || session_id.is_empty() {
return None;
}
Some(ProviderSessionRef {
harness: harness.to_string(),
session_id: session_id.to_string(),
})
}
fn body_has_harness(body: &BodyProvenance) -> bool {
body.harness
.as_deref()
.is_some_and(|harness| !harness.trim().is_empty())
}
fn claim_answers(inner: &mut Inner, answers: Vec<String>) -> Vec<MessageId> {
let mut valid = Vec::new();
for id in answers {
let id = MessageId(id);
if let Some(pos) = inner.pending_messages.iter().position(|m| m.id == id) {
inner.pending_messages.remove(pos);
valid.push(id);
} else {
tracing::warn!(
id = %id,
"resident answered an unknown or already-consumed message; dropped"
);
}
}
valid
}
fn snapshot_tail_locked(inner: &Inner, limit: Option<usize>) -> Vec<ChatTurn> {
let open_count = usize::from(inner.open.is_some());
let total = inner.thread.len() + open_count;
let take = limit.unwrap_or(total).min(total);
let take_open = take.min(open_count);
let take_thread = take - take_open;
let mut turns = inner.thread[inner.thread.len() - take_thread..].to_vec();
if take_open == 1 {
turns.extend(inner.open.as_ref().map(|open| open.turn.clone()));
}
turns
}
fn add_opt(a: Option<u64>, b: Option<u64>) -> Option<u64> {
match (a, b) {
(None, None) => None,
(a, b) => Some(a.unwrap_or(0) + b.unwrap_or(0)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::Path;
fn turn_seq(id: &str) -> u64 {
id.strip_prefix("turn-")
.and_then(|n| n.parse().ok())
.expect("turn id minted from journal seq")
}
fn msg_id(turn: &ChatTurn) -> String {
format!("msg-{}", turn_seq(&turn.id))
}
fn progress_turn(text: &str) -> ChatTurn {
ChatTurn {
id: String::new(),
role: ChatRole::Assistant,
text: text.to_string(),
status: Lifecycle::Completed,
items: Vec::new(),
created_at: String::new(),
from: None,
body: None,
activity: None,
}
}
fn open_runtime(repo: &Path) -> Arc<WaveRuntime> {
WaveRuntime::open("ship".into(), repo.to_path_buf()).expect("open runtime")
}
fn d_opened(answers: &[&str]) -> ResidentDelta {
ResidentDelta::TurnOpened {
answers: answers.iter().map(|s| s.to_string()).collect(),
}
}
fn d_text(text: &str) -> ResidentDelta {
ResidentDelta::TurnText { text: text.into() }
}
fn d_tool() -> ResidentDelta {
ResidentDelta::TurnItem {
item: ConversationItem::Tool {
id: "item-tool".into(),
name: "Bash".into(),
status: Lifecycle::Completed,
input: None,
output: Some("cargo test".into()),
},
}
}
fn d_thought(id: &str, text: &str) -> ResidentDelta {
ResidentDelta::TurnItem {
item: ConversationItem::Thought {
id: id.into(),
text: text.into(),
},
}
}
fn d_usage(input: u64, output: u64) -> ResidentDelta {
ResidentDelta::TurnUsage {
input_tokens: Some(input),
output_tokens: Some(output),
cache_read_tokens: None,
}
}
fn d_finished(status: Lifecycle) -> ResidentDelta {
ResidentDelta::TurnFinished {
status,
cost_usd: None,
reason: None,
}
}
fn complete_body(rt: &WaveRuntime, harness: &str, session_id: Option<&str>) {
let step = rt
.ensure_playhead()
.expect("initialize playhead")
.now
.expect("wave has a current step");
let mut body = BodyProvenance::for_step(&step, rt.repo_root());
body.harness = Some(harness.to_string());
body.session_id = session_id.map(str::to_string);
let body_id = body.body_id.clone();
rt.apply_resident_delta(ResidentDelta::BodyStarted { body });
rt.apply_resident_delta(d_opened(&[]));
rt.apply_resident_delta(d_text("done"));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
rt.apply_resident_delta(ResidentDelta::BodyFinished {
body_id,
outcome: StepOutcome::Completed,
reason: "completed".to_string(),
});
}
#[test]
fn turns_get_monotonic_ids_from_the_journal() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let a = rt.append_finalized_turn(progress_turn("one"), Vec::new());
let b = rt.append_finalized_turn(progress_turn("two"), Vec::new());
assert!(turn_seq(&b.id) > turn_seq(&a.id));
assert_eq!(rt.thread_snapshot().len(), 2);
}
#[test]
fn narrated_turns_no_longer_blob_memory() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
rt.append_finalized_turn(progress_turn("landed the parser"), Vec::new());
assert_eq!(rt.memory().read(), "");
}
#[test]
fn deliver_appends_user_turn_and_broadcasts() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut rx = rt.subscribe_inbox();
let turn = rt
.deliver(MessageOp::Message, "how goes it?".into())
.expect("user turn");
assert_eq!(turn.role, ChatRole::User);
assert_eq!(turn.text, "how goes it?");
let InboxItem::Message(msg) = rx.try_recv().expect("inbox message") else {
panic!("expected a message inbox item");
};
assert_eq!(msg.text, "how goes it?");
assert_eq!(msg.op, MessageOp::Message);
assert_eq!(msg.id, MessageId(msg_id(&turn)));
assert_eq!(rt.pending_messages().len(), 1);
}
#[test]
fn task_observation_is_typed_idempotent_and_replayable() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut rx = rt.subscribe_inbox();
let observation = crate::task::TaskObservation {
session_id: crate::task::TaskSessionId::from_raw("ts_example"),
issue_identifier: "INF-123".to_string(),
event_id: 7,
control_source: None,
event: crate::task::TaskEventKind::DecisionRequested {
decision_id: crate::child_session::ChildDecisionId::new(),
prompt: "Approve the plan?".to_string(),
options: vec!["approve".to_string(), "revise".to_string()],
},
};
assert!(rt.deliver_task_observation(observation.clone()));
assert!(!rt.deliver_task_observation(observation.clone()));
assert!(matches!(
rx.try_recv().expect("live Task observation"),
InboxItem::Task(ref received) if received == &observation
));
let sub = rt.subscribe_with_snapshot(None);
assert_eq!(sub.pending.len(), 1);
assert_eq!(
sub.tasks.get(&MessageId(observation.inbox_id())),
Some(&observation)
);
drop(rt);
let replayed = open_runtime(tmp.path()).subscribe_with_snapshot(None);
assert_eq!(replayed.pending.len(), 1);
assert_eq!(
replayed.tasks.get(&MessageId(observation.inbox_id())),
Some(&observation)
);
}
#[test]
fn bounded_subscription_tails_replay_but_keeps_live_turns() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
for i in 0..5 {
rt.deliver(MessageOp::Message, format!("message {i}"))
.expect("user turn");
}
let mut sub = rt.subscribe_with_snapshot(Some(2));
assert_eq!(
sub.turns
.iter()
.map(|turn| turn.text.as_str())
.collect::<Vec<_>>(),
vec!["message 3", "message 4"]
);
rt.deliver(MessageOp::Message, "message 5".into())
.expect("live user turn");
assert_eq!(
sub.turn_rx
.try_recv()
.expect("live frame")
.expect_whole()
.text,
"message 5",
"the limit applies only to replay"
);
}
#[test]
fn project_observation_is_typed_idempotent_and_replayable() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut rx = rt.subscribe_inbox();
let observation = crate::project_session::ProjectObservation {
session_id: crate::project_session::ProjectSessionId::from_raw("ps_example"),
project: "developer-efficiency".to_string(),
event_id: 8,
control_source: None,
event: crate::project_session::ProjectEventKind::Completed {
summary: "all KRs hold".to_string(),
},
};
assert!(rt.deliver_project_observation(observation.clone()));
assert!(!rt.deliver_project_observation(observation.clone()));
assert!(matches!(
rx.try_recv().expect("live Project observation"),
InboxItem::Project(ref received) if received == &observation
));
let sub = rt.subscribe_with_snapshot(None);
assert_eq!(
sub.projects.get(&MessageId(observation.inbox_id())),
Some(&observation)
);
drop(rt);
let replayed = open_runtime(tmp.path()).subscribe_with_snapshot(None);
assert_eq!(
replayed.projects.get(&MessageId(observation.inbox_id())),
Some(&observation)
);
}
#[test]
fn deliver_say_journals_attribution_and_queues_for_the_loop() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut rx = rt.subscribe_inbox();
let from = "worker".to_string();
let turn = rt.deliver_say("PR landed; one surprise in the fold".into(), from.clone());
assert_eq!(turn.role, ChatRole::User);
assert_eq!(turn.from.as_deref(), Some("worker"));
let InboxItem::Message(msg) = rx.try_recv().expect("inbox item") else {
panic!("expected a message inbox item");
};
assert_eq!(msg.op, MessageOp::Say);
assert_eq!(msg.from, Some(from.clone()));
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let EventKind::UserMessage {
op, from: stored, ..
} = &events[0].kind
else {
panic!("expected UserMessage");
};
assert_eq!(*op, MessageOp::Say);
assert_eq!(stored.as_ref(), Some(&from));
let rt2 = open_runtime(tmp.path());
let pending = rt2.pending_messages();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].from, Some(from));
assert_eq!(rt2.thread_snapshot()[0].from.as_deref(), Some("worker"));
}
#[test]
fn update_memory_writes_the_origin_file_and_add_publishes_delta() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
rt.update_memory("# Ship\n\n- fold is truth\n", "fold is truth")
.expect("update");
rt.append_memory("fold is truth", vec![]).expect("append");
rt.append_memory("bullets append", vec![]).expect("append");
assert_eq!(
rt.memory().read(),
"# Ship\n\n- fold is truth\n",
"adds do not accrete raw facts into the compiled ORIGIN file"
);
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let summaries: Vec<&str> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::MemoryUpdated { summary } => Some(summary.as_str()),
_ => None,
})
.collect();
let facts: Vec<&str> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::MemoryAdded { fact, .. } => Some(fact.as_str()),
_ => None,
})
.collect();
assert_eq!(summaries, vec!["fold is truth"]);
assert_eq!(facts, vec!["fold is truth", "bullets append"]);
}
#[test]
fn subscription_replays_full_memory_facts_in_order() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let long_fact = "workers report via lf radio pub with the full useful detail";
rt.append_memory(long_fact, vec![]).expect("append");
rt.append_memory("second fact", vec![]).expect("append");
let sub = rt.subscribe_with_snapshot(None);
assert_eq!(
sub.memory_adds,
vec![long_fact.to_string(), "second fact".to_string()]
);
assert!(
sub.memory_add_rx.is_empty(),
"snapshot facts do not replay live"
);
}
#[test]
fn memory_add_replay_buffer_rebuilds_from_journal() {
let tmp = tempfile::tempdir().expect("tempdir");
{
let rt = open_runtime(tmp.path());
rt.append_memory("first", vec![]).expect("append");
rt.append_memory("second", vec![]).expect("append");
rt.update_memory("# Ship\n\ncompiled\n", "compiled")
.expect("update");
rt.append_memory("third", vec![]).expect("append");
}
let rt = open_runtime(tmp.path());
let sub = rt.subscribe_with_snapshot(None);
assert_eq!(
sub.memory_adds,
vec!["third".to_string()],
"the replay buffer rebuilds adds since the last externalization"
);
}
#[test]
fn deliver_interrupt_is_a_control_item_not_a_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut rx = rt.subscribe_inbox();
rt.deliver_interrupt();
assert!(matches!(
rx.try_recv().expect("inbox item"),
InboxItem::Interrupt
));
assert!(rt.thread_snapshot().is_empty());
assert!(rt.pending_messages().is_empty());
}
#[test]
fn illegal_transition_is_refused_and_leaves_state_alone() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
assert_eq!(rt.loop_state(), LoopState::Idle);
assert!(!rt.transition(
LoopState::Interrupting {
turn_id: "turn-1".into()
},
"test"
));
assert_eq!(rt.loop_state(), LoopState::Idle);
assert!(rt.transition(
LoopState::Turning {
turn_id: "turn-1".into()
},
"test"
));
assert!(!rt.transition(
LoopState::Turning {
turn_id: "turn-2".into()
},
"test"
));
assert!(rt.transition(LoopState::Idle, "test"));
}
#[test]
fn resident_deltas_journal_a_turn_and_commit_it() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
rt.apply_resident_delta(d_opened(&[]));
assert_eq!(
rt.loop_state().name(),
"turning",
"mid-turn the loop is Turning"
);
rt.apply_resident_delta(d_text("hello"));
rt.apply_resident_delta(d_tool());
rt.apply_resident_delta(d_usage(10, 4));
rt.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: Some(0.02),
reason: None,
});
assert_eq!(rt.loop_state(), LoopState::Idle, "back to idle after turn");
let thread = rt.thread_snapshot();
assert_eq!(thread.len(), 1);
let turn = &thread[0];
assert_eq!(turn.text, "hello");
assert_eq!(turn.items.len(), 1);
assert_eq!(turn.status, Lifecycle::Completed);
turn_seq(&turn.id);
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let usage = events
.iter()
.find_map(|e| match &e.kind {
EventKind::TurnFinished { usage, .. } => Some(usage.clone()),
_ => None,
})
.expect("TurnFinished journaled");
assert_eq!(usage.input_tokens, Some(10));
assert_eq!(usage.output_tokens, Some(4));
assert_eq!(usage.cost_usd, Some(0.02));
}
#[test]
fn provider_session_survives_runtime_reopen() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
complete_body(&rt, "codex", Some("thread-resume"));
assert_eq!(
rt.latest_provider_session(),
Some(ProviderSessionRef {
harness: "codex".to_string(),
session_id: "thread-resume".to_string(),
})
);
drop(rt);
let reopened = open_runtime(tmp.path());
assert_eq!(
reopened.latest_provider_session(),
Some(ProviderSessionRef {
harness: "codex".to_string(),
session_id: "thread-resume".to_string(),
})
);
}
#[test]
fn newest_harness_body_masks_an_older_provider_session() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
complete_body(&rt, "codex", Some("codex-thread"));
complete_body(&rt, "claude", None);
assert_eq!(rt.latest_provider_session(), None);
}
#[test]
fn turn_opened_answers_are_validated_against_the_pending_fold() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let m1 = rt
.deliver(MessageOp::Message, "first".into())
.expect("user turn");
let m2 = rt
.deliver(MessageOp::Message, "second".into())
.expect("user turn");
assert_eq!(rt.pending_messages().len(), 2);
rt.apply_resident_delta(d_opened(&[&msg_id(&m1), &msg_id(&m2), "msg-999"]));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
assert!(
rt.pending_messages().is_empty(),
"claimed messages leave the live pending fold"
);
rt.apply_resident_delta(d_opened(&[&msg_id(&m1)]));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let answers: Vec<Vec<MessageId>> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::TurnStarted { answers, .. } => Some(answers.clone()),
_ => None,
})
.collect();
assert_eq!(answers.len(), 2);
assert_eq!(
answers[0],
vec![MessageId(msg_id(&m1)), MessageId(msg_id(&m2))],
"valid ids journaled, the ghost dropped"
);
assert!(answers[1].is_empty(), "already-consumed ids never re-claim");
let fold = fold_thread(&events);
assert!(fold.pending_messages.is_empty());
}
#[test]
fn open_turn_streams_growing_snapshots_then_the_terminal_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let sub = rt.subscribe_with_snapshot(None);
assert!(sub.turns.is_empty());
assert_eq!(sub.state, LoopState::Idle);
let mut frames = sub.turn_rx;
let mut states = sub.state_rx;
rt.apply_resident_delta(d_opened(&[]));
let opened = frames
.try_recv()
.expect("opened frame")
.expect_whole()
.clone();
assert_eq!(opened.status, Lifecycle::Running);
assert_eq!(opened.text, "");
let mut reconstruction = opened.clone();
rt.apply_resident_delta(d_text("thinking"));
let text_delta = frames
.try_recv()
.expect("text delta")
.expect_delta()
.clone();
assert_eq!(text_delta.turn_id, opened.id);
reconstruction.absorb_item(text_delta.item);
assert_eq!(reconstruction.text, "thinking");
let mid = rt.thread_snapshot();
assert_eq!(mid.len(), 1);
assert_eq!(mid[0].id, opened.id);
assert_eq!(mid[0].status, Lifecycle::Running);
assert_eq!(mid[0].text, "thinking");
rt.apply_resident_delta(d_tool());
let item_delta = frames
.try_recv()
.expect("item delta")
.expect_delta()
.clone();
assert_eq!(item_delta.turn_id, opened.id);
reconstruction.absorb_item(item_delta.item);
assert_eq!(reconstruction.items.len(), 1);
assert_eq!(reconstruction.text, "thinking");
rt.apply_resident_delta(d_usage(10, 5));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let terminal = frames
.try_recv()
.expect("terminal frame")
.expect_whole()
.clone();
assert_eq!(terminal.id, opened.id);
assert_eq!(terminal.status, Lifecycle::Completed);
assert_eq!(reconstruction.text, terminal.text);
assert_eq!(reconstruction.items, terminal.items);
let after = rt.thread_snapshot();
assert_eq!(after.len(), 1);
assert_eq!(after[0].status, Lifecycle::Completed);
assert!(frames.try_recv().is_err(), "no extra frames");
assert!(matches!(
states.try_recv().expect("turning state frame"),
LoopState::Turning { .. }
));
assert_eq!(
states.try_recv().expect("idle state frame"),
LoopState::Idle
);
assert!(states.try_recv().is_err(), "no extra state frames");
}
#[test]
fn a_text_delta_never_carries_the_turns_accumulated_items() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut frames = rt.subscribe_with_snapshot(None).turn_rx;
rt.apply_resident_delta(d_opened(&[]));
let _ = frames.try_recv().expect("opened frame");
let big_output = "x".repeat(100_000);
rt.apply_resident_delta(ResidentDelta::TurnItem {
item: ConversationItem::Tool {
id: "item-big".into(),
name: "Bash".into(),
status: Lifecycle::Completed,
input: None,
output: Some(big_output.clone()),
},
});
let _ = frames.try_recv().expect("tool delta");
rt.apply_resident_delta(d_text("tiny"));
let TurnBroadcast::Delta(frame) = frames.try_recv().expect("text delta") else {
panic!("a text increment must broadcast a delta, not a whole turn");
};
assert!(
!frame.json.contains(&big_output),
"the delta frame must not re-send the turn's accumulated items"
);
assert!(
frame.json.len() < 200,
"a token's frame stays tiny ({} bytes) regardless of turn size",
frame.json.len()
);
}
#[test]
fn empty_thoughts_never_enter_the_thread_or_journal() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut frames = rt.subscribe_with_snapshot(None).turn_rx;
rt.apply_resident_delta(d_opened(&[]));
let opened = frames.try_recv().expect("opened frame");
assert!(opened.expect_whole().items.is_empty());
rt.apply_resident_delta(d_thought("empty", " \n\t"));
assert!(frames.try_recv().is_err(), "empty thought emits no frame");
assert!(rt.thread_snapshot()[0].items.is_empty());
rt.apply_resident_delta(d_thought("real", "checking the retry"));
let grown = frames.try_recv().expect("real thought frame");
assert!(matches!(
&grown.expect_delta().item,
ConversationItem::Thought { id, text }
if id == "real" && text == "checking the retry"
));
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let thought_items: Vec<&ConversationItem> = events
.iter()
.filter_map(|event| match &event.kind {
EventKind::TurnItem { item, .. } => Some(item),
_ => None,
})
.collect();
assert_eq!(thought_items.len(), 1, "only the real thought is journaled");
}
#[test]
fn force_finalize_closes_the_turn_and_drops_late_deltas() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let step = rt
.ensure_playhead()
.expect("initialize playhead")
.now
.expect("wave has a current step");
let body = BodyProvenance {
body_id: "body-dead".into(),
invocation_id: step.invocation_id,
step_index: step.index,
flow: step.flow,
step: step.step,
iteration: step.iteration,
session_id: Some("session-dead".into()),
harness: Some("codex".into()),
model: None,
host: "host".into(),
worktree: tmp.path().display().to_string(),
started_at: "2026-07-09T00:00:00Z".into(),
ended_at: None,
termination_reason: None,
};
rt.apply_resident_delta(ResidentDelta::BodyStarted { body });
rt.apply_resident_delta(d_opened(&[]));
rt.apply_resident_delta(d_text("half"));
rt.apply_resident_delta(ResidentDelta::LoopState {
to: ResidentStateTo::Interrupting,
reason: "user interrupt".into(),
});
assert_eq!(rt.loop_state().name(), "interrupting");
assert!(rt.force_finalize_open_turn(Lifecycle::Interrupted, "deadline"));
assert_eq!(rt.loop_state(), LoopState::Idle);
let thread = rt.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].status, Lifecycle::Interrupted);
assert_eq!(thread[0].text, "half");
let body = thread[0].body.as_ref().expect("turn keeps body provenance");
assert!(body.ended_at.is_some());
assert_eq!(body.termination_reason.as_deref(), Some("deadline"));
let playhead = rt.playhead().expect("playhead survives finalization");
assert!(playhead.active.is_none(), "the dead body releases its seat");
assert_eq!(
playhead.now.expect("failed step remains selected").index,
0,
"an interrupted body retries the same logical step"
);
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let fold = crate::wave::journal::fold_thread(&events);
assert!(fold.open.is_empty());
assert_eq!(fold.turns.last().unwrap().status, Lifecycle::Interrupted);
assert!(
fold.playhead.expect("replayed playhead").active.is_none(),
"replay releases the dead body too"
);
assert!(!rt.force_finalize_open_turn(Lifecycle::Interrupted, "again"));
let journal_len = events.len();
rt.apply_resident_delta(d_text("late text"));
rt.apply_resident_delta(d_usage(1, 1));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
assert_eq!(events.len(), journal_len, "late deltas journal nothing");
assert_eq!(rt.thread_snapshot().len(), 1, "thread untouched");
assert_eq!(rt.loop_state(), LoopState::Idle, "no double transition");
rt.apply_resident_delta(d_opened(&[]));
rt.apply_resident_delta(d_text("fresh"));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let thread = rt.thread_snapshot();
assert_eq!(thread.len(), 2);
assert_eq!(thread[1].text, "fresh");
assert_eq!(thread[1].status, Lifecycle::Completed);
}
#[test]
fn channel_role_compares_against_the_sanitized_wave_name() {
assert_eq!(channel_role("web/ui", "web/ui"), Some(ChannelRole::Primary));
assert_eq!(channel_role("web/ui", "web-ui"), Some(ChannelRole::Primary));
assert_eq!(
channel_role("web/ui", "web-ui.148e0e02"),
Some(ChannelRole::Child)
);
assert_eq!(channel_role("web/ui", "web/ui.148e"), None);
assert_eq!(channel_role("goals", "goals"), Some(ChannelRole::Primary));
assert_eq!(
channel_role("goals", "goals.148e0e02"),
Some(ChannelRole::Child)
);
assert_eq!(channel_role("goals", "goalsmith"), None);
assert_eq!(channel_role("goals", "concerto"), None);
}
#[test]
fn a_sanitized_wave_recognizes_its_family_by_the_sanitized_name() {
let tmp = tempfile::tempdir().expect("tempdir");
let origin = tmp.path().join("repo");
let rt = WaveRuntime::open("web/ui".into(), origin).expect("open runtime");
assert_eq!(rt.channel_name(), "web-ui");
assert!(rt.is_primary("web-ui"));
assert!(
rt.is_primary("web/ui"),
"the raw spelling still addresses it"
);
assert!(rt.in_family("web-ui.148e"));
assert!(!rt.in_family("web-ui-other"));
}
#[test]
fn failed_turn_requeues_its_claimed_messages() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let m1 = rt
.deliver(MessageOp::Message, "do the thing".into())
.expect("user turn");
rt.apply_resident_delta(d_opened(&[&msg_id(&m1)]));
assert!(rt.pending_messages().is_empty(), "claimed at open");
rt.apply_resident_delta(d_finished(Lifecycle::Failed));
let pending = rt.pending_messages();
assert_eq!(pending.len(), 1, "failed turn requeues its claim");
assert_eq!(pending[0].text, "do the thing");
let rt2 = open_runtime(tmp.path());
assert_eq!(rt2.pending_messages().len(), 1);
let tmp3 = tempfile::tempdir().expect("tempdir");
let rt3 = open_runtime(tmp3.path());
let m2 = rt3
.deliver(MessageOp::Message, "second".into())
.expect("user turn");
rt3.apply_resident_delta(d_opened(&[&msg_id(&m2)]));
rt3.apply_resident_delta(d_finished(Lifecycle::Completed));
assert!(
rt3.pending_messages().is_empty(),
"a completed turn consumes its claim for good"
);
}
#[test]
fn boot_janitor_requeues_a_crashed_turns_claims() {
let tmp = tempfile::tempdir().expect("tempdir");
let claimed = {
let rt = open_runtime(tmp.path());
let m = rt
.deliver(MessageOp::Message, "answer me".into())
.expect("user turn");
rt.apply_resident_delta(d_opened(&[&msg_id(&m)]));
assert!(rt.pending_messages().is_empty());
msg_id(&m)
};
let rt = open_runtime(tmp.path());
let pending = rt.pending_messages();
assert_eq!(pending.len(), 1, "crashed turn's claim is requeued");
assert_eq!(pending[0].id, MessageId(claimed));
let rt2 = open_runtime(tmp.path());
assert_eq!(rt2.pending_messages().len(), 1);
}
#[test]
fn resident_requeue_undoes_a_claim_at_most_once() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let m = rt
.deliver(MessageOp::Steer, "steer".into())
.expect("user turn");
rt.apply_resident_delta(d_opened(&[&msg_id(&m)]));
rt.apply_resident_delta(ResidentDelta::MessagesRequeued {
ids: vec![msg_id(&m)],
});
assert_eq!(
rt.pending_messages().len(),
1,
"undone claim back to pending"
);
rt.apply_resident_delta(d_finished(Lifecycle::Failed));
assert_eq!(rt.pending_messages().len(), 1, "no double requeue");
}
#[test]
fn paused_wave_refuses_to_start_turns() {
let tmp = tempfile::tempdir().expect("tempdir");
let origin = tmp.path();
std::fs::create_dir_all(origin.join("wave/ship")).unwrap();
std::fs::write(
origin.join("wave/ship/GOAL.md"),
"---\npaused: true\n---\nShip it.\n",
)
.unwrap();
let rt = open_runtime(origin);
assert!(rt.paused(), "GOAL.md says paused");
let m = rt
.deliver(MessageOp::Message, "go".into())
.expect("user turn");
rt.apply_resident_delta(d_opened(&[&msg_id(&m)]));
rt.apply_resident_delta(d_text("working"));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
assert!(
rt.thread_snapshot()
.iter()
.all(|t| t.role == ChatRole::User),
"paused: no assistant turn started"
);
assert_eq!(rt.pending_messages().len(), 1, "the message waits");
std::fs::write(
origin.join("wave/ship/GOAL.md"),
"---\npaused: false\n---\nShip it.\n",
)
.unwrap();
assert!(!rt.paused());
rt.apply_resident_delta(d_opened(&[&msg_id(&m)]));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
assert!(
rt.pending_messages().is_empty(),
"unpaused turn answered it"
);
}
#[test]
fn turn_steered_consumes_against_live_or_just_closed_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let m1 = rt
.deliver(MessageOp::Steer, "steer me".into())
.expect("user turn");
let m2 = rt
.deliver(MessageOp::Steer, "me too".into())
.expect("user turn");
rt.apply_resident_delta(ResidentDelta::TurnSteered {
answers: vec![msg_id(&m1)],
});
assert_eq!(rt.pending_messages().len(), 2, "kept pending");
rt.apply_resident_delta(d_opened(&[]));
rt.apply_resident_delta(ResidentDelta::TurnSteered {
answers: vec![msg_id(&m1)],
});
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
rt.apply_resident_delta(ResidentDelta::TurnSteered {
answers: vec![msg_id(&m2)],
});
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let assistant_turn = rt
.thread_snapshot()
.iter()
.find(|turn| turn.role == ChatRole::Assistant)
.map(|turn| turn.id.clone())
.expect("assistant turn");
let steered: Vec<_> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::TurnSteered { turn_id, answers } => {
Some((turn_id.clone(), answers.clone()))
}
_ => None,
})
.collect();
assert_eq!(
steered,
vec![
(assistant_turn.clone(), vec![MessageId(msg_id(&m1))]),
(assistant_turn, vec![MessageId(msg_id(&m2))]),
],
"both markers name the turn that heard the text"
);
let fold = fold_thread(&events);
assert!(
fold.pending_messages.is_empty(),
"consumed messages never re-send: {:?}",
fold.pending_messages
);
}
#[test]
fn turn_steered_fallback_survives_restart() {
let tmp = tempfile::tempdir().expect("tempdir");
let assistant_id = {
let rt = open_runtime(tmp.path());
rt.apply_resident_delta(d_opened(&[]));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
rt.thread_snapshot()
.iter()
.find(|turn| turn.role == ChatRole::Assistant)
.expect("assistant turn")
.id
.clone()
};
let rt = open_runtime(tmp.path());
let steer = rt
.deliver(MessageOp::Steer, "steer me".into())
.expect("user turn");
rt.apply_resident_delta(ResidentDelta::TurnSteered {
answers: vec![msg_id(&steer)],
});
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let steered_turn = events
.iter()
.find_map(|e| match &e.kind {
EventKind::TurnSteered { turn_id, .. } => Some(turn_id.clone()),
_ => None,
})
.expect("TurnSteered journaled");
assert_eq!(
steered_turn, assistant_id,
"fallback names the first life's assistant turn, not the user turn"
);
}
#[test]
fn only_the_served_mind_journals() {
let tmp = tempfile::tempdir().expect("tempdir");
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let rt = WaveRuntime::open("ship".into(), origin.clone()).expect("open runtime");
rt.deliver(MessageOp::Message, "to the wave".into())
.expect("user turn");
rt.deliver_say("to a".into(), "ship.a".into());
let wave = rt.thread_snapshot();
assert_eq!(wave.len(), 2);
assert_eq!(wave[0].text, "to the wave");
assert_eq!(wave[1].from.as_deref(), Some("ship.a"));
assert_eq!(
loopflow_test_support::journal_files_under(tmp.path()),
vec![journal_path(&origin, "ship")]
);
let events = crate::wave::journal::read_events(&journal_path(&origin, "ship"));
let messages = events
.iter()
.filter(|e| matches!(e.kind, EventKind::UserMessage { .. }))
.count();
assert_eq!(messages, 2, "the wave's message and the folded report");
let rt2 = WaveRuntime::open("ship".into(), origin.clone()).expect("reopen");
let pending = rt2.pending_messages();
let texts: Vec<_> = pending.iter().map(|m| m.text.as_str()).collect();
assert_eq!(texts.len(), 2);
assert!(texts.contains(&"to the wave") && texts.contains(&"to a"));
}
#[test]
fn subscription_carries_pending_replay_and_live_inbox() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
rt.deliver(MessageOp::Message, "before".into())
.expect("user turn");
let mut sub = rt.subscribe_with_snapshot(None);
assert_eq!(sub.pending.len(), 1);
assert_eq!(sub.pending[0].text, "before");
assert!(sub.inbox_rx.try_recv().is_err(), "no frames from before");
rt.deliver(MessageOp::Message, "after".into())
.expect("user turn");
rt.deliver_interrupt();
let InboxItem::Message(live) = sub.inbox_rx.try_recv().expect("live frame") else {
panic!("expected message");
};
assert_eq!(live.text, "after");
assert!(matches!(
sub.inbox_rx.try_recv().expect("interrupt frame"),
InboxItem::Interrupt
));
}
#[test]
fn concurrent_deliveries_keep_inbox_order_equal_to_journal_order() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let mut rx = rt.subscribe_inbox();
let mut handles = Vec::new();
for writer in 0..4 {
let rt = rt.clone();
handles.push(std::thread::spawn(move || {
for i in 0..50 {
rt.deliver(MessageOp::Message, format!("m-{writer}-{i}"))
.expect("user turn");
}
}));
}
for handle in handles {
handle.join().expect("writer thread");
}
let mut inbox_ids = Vec::new();
while let Ok(item) = rx.try_recv() {
let InboxItem::Message(message) = item else {
panic!("only messages were delivered");
};
inbox_ids.push(message.id);
}
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let journal_ids: Vec<MessageId> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::UserMessage { id, .. } => Some(id.clone()),
_ => None,
})
.collect();
assert_eq!(inbox_ids.len(), 200);
assert_eq!(
inbox_ids, journal_ids,
"inbox consumption order == journal fold order"
);
}
}