use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use sha2::{Digest, Sha256};
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::ProjectObservation;
use crate::task::TaskObservation;
use crate::wave::chat::{ChatBacking, ChatMessageSource, ConversationEpoch, WaveChatMessage};
#[cfg(test)]
use crate::wave::journal::JournalAppendStage;
use crate::wave::journal::{
fold_thread, journal_path, project_observation_message, promotion_wake_message,
restore_pending, task_observation_message, ConversationEpochImport, DiscordAttachment,
DiscordChatBinding, DiscordDelivery, DiscordMessagePart, DiscordMessageSource, EventKind,
Journal, JournalAppendError, MessageId, MessageOp, PendingMessage, Usage,
};
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};
use crate::wave::PromotionWake;
const TURN_BROADCAST_CAPACITY: usize = 256;
const STATE_BROADCAST_CAPACITY: usize = 64;
const PLAYHEAD_BROADCAST_CAPACITY: usize = 64;
const INBOX_BROADCAST_CAPACITY: usize = 256;
#[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),
Promotion {
parent_wave_id: crate::id::WaveId,
parent: String,
},
Interrupt,
Skip,
}
#[derive(Debug)]
pub struct Subscription {
pub epoch: ConversationEpoch,
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 pending: Vec<PendingMessage>,
pub tasks: HashMap<MessageId, TaskObservation>,
pub projects: HashMap<MessageId, ProjectObservation>,
pub(crate) promotions: HashMap<MessageId, PromotionWake>,
pub inbox_rx: broadcast::Receiver<InboxItem>,
}
#[derive(Debug, Clone)]
pub struct DiscordSnapshot {
pub attachment: Option<DiscordAttachment>,
pub deliveries: Vec<DiscordDelivery>,
}
#[derive(Debug, Clone, Copy)]
enum DiscordInput {
Provider,
Authored(MessageOp),
}
impl DiscordInput {
fn op(self) -> MessageOp {
match self {
Self::Provider => MessageOp::Message,
Self::Authored(op) => op,
}
}
}
#[derive(Debug)]
struct OpenTurn {
turn: ChatTurn,
usage: Usage,
text_items: usize,
claims: Vec<MessageId>,
}
#[derive(Debug)]
struct Inner {
journal: Journal,
thread: Vec<ChatTurn>,
conversation_epochs: Vec<ConversationEpoch>,
conversation_epoch_turns: HashMap<String, Vec<String>>,
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>,
discord: Option<DiscordAttachment>,
discord_deliveries: HashMap<String, DiscordDelivery>,
tasks: HashMap<MessageId, TaskObservation>,
projects: HashMap<MessageId, ProjectObservation>,
promotions: HashMap<MessageId, PromotionWake>,
}
#[derive(Debug)]
pub struct WaveRuntime {
name: String,
repo_root: PathBuf,
inner: Mutex<Inner>,
turn_tx: broadcast::Sender<TurnBroadcast>,
state_tx: broadcast::Sender<LoopState>,
playhead_tx: broadcast::Sender<PlayheadView>,
inbox_tx: broadcast::Sender<InboxItem>,
resident_expected: AtomicBool,
}
#[derive(Debug, thiserror::Error)]
pub enum ChatWriteError {
#[error("this Wave chat is backed by Discord")]
OpenDiscord,
#[error(transparent)]
Journal(#[from] JournalAppendError),
}
impl WaveRuntime {
pub fn open(name: String, repo_root: PathBuf) -> anyhow::Result<Arc<Self>> {
Self::open_with_backing(name, repo_root, ChatBacking::Local)
}
pub fn open_with_backing(
name: String,
repo_root: PathBuf,
backing: ChatBacking,
) -> 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
};
initialize_conversation_epoch(
&mut journal,
&mut fold.conversation_epochs,
&mut fold.conversation_epoch_turns,
&fold.turns,
&fold.discord_turn_bindings,
backing.clone(),
);
let planned_turns = fold
.discord_deliveries
.values()
.map(|delivery| delivery.turn_id.clone())
.collect::<std::collections::HashSet<_>>();
let active_epoch = fold
.conversation_epochs
.last()
.cloned()
.expect("conversation epoch initialized before delivery recovery");
for (turn_id, claims) in &fold.completed_claims {
if planned_turns.contains(turn_id) {
continue;
}
let Some(turn) = fold.turns.iter().find(|turn| &turn.id == turn_id) else {
continue;
};
if turn_journal_seq(turn).is_none_or(|seq| seq <= active_epoch.journal_seq) {
continue;
}
let Some(binding) = active_epoch.backing.discord_binding() else {
continue;
};
let Some(delivery) = build_discord_delivery(turn, claims, &fold.messages, &binding)
else {
continue;
};
journal.append(|_| EventKind::DiscordChatSendPlanned {
delivery_id: delivery.delivery_id.clone(),
turn_id: delivery.turn_id.clone(),
binding: Some(delivery.binding.clone()),
sources: delivery.sources.clone(),
parts: delivery.parts.clone(),
});
fold.discord_deliveries
.insert(delivery.delivery_id.clone(), delivery);
}
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 (inbox_tx, _) = broadcast::channel(INBOX_BROADCAST_CAPACITY);
Ok(Arc::new(Self {
name,
repo_root,
inner: Mutex::new(Inner {
journal,
thread: fold.turns,
conversation_epochs: fold.conversation_epochs,
conversation_epoch_turns: fold.conversation_epoch_turns,
open: None,
drop_deltas_until_opened: false,
state,
playhead,
last_assistant_turn_id,
pending_messages: fold.pending_messages,
messages: fold.messages,
discord: fold.discord,
discord_deliveries: fold.discord_deliveries,
tasks: fold.tasks,
projects: fold.projects,
promotions: fold.promotions,
}),
turn_tx,
state_tx,
playhead_tx,
inbox_tx,
resident_expected: AtomicBool::new(false),
}))
}
pub fn name(&self) -> &str {
&self.name
}
pub fn repo_root(&self) -> &std::path::Path {
&self.repo_root
}
pub fn active_conversation_epoch(&self) -> ConversationEpoch {
self.inner()
.conversation_epochs
.last()
.cloned()
.expect("an open runtime always has an active conversation epoch")
}
pub fn conversation_epochs(&self) -> Vec<ConversationEpoch> {
self.inner().conversation_epochs.clone()
}
pub fn is_imported_conversation_epoch(&self, epoch_id: &str) -> bool {
self.inner().conversation_epoch_turns.contains_key(epoch_id)
}
pub fn chat_messages(
&self,
epoch_id: Option<&str>,
limit: Option<usize>,
) -> Vec<WaveChatMessage> {
let inner = self.inner();
let selected = match epoch_id {
Some(id) => inner
.conversation_epochs
.iter()
.find(|epoch| epoch.id == id),
None => inner.conversation_epochs.last(),
};
let Some(epoch) = selected else {
return Vec::new();
};
if !matches!(epoch.backing, ChatBacking::Local) {
return Vec::new();
}
let imported_turns = inner.conversation_epoch_turns.get(&epoch.id);
let end_seq = inner
.conversation_epochs
.iter()
.find(|candidate| candidate.number == epoch.number + 1)
.map(|candidate| candidate.journal_seq)
.unwrap_or(u64::MAX);
let turns = snapshot_tail_locked(&inner, None)
.into_iter()
.filter_map(|turn| {
let journal_seq = turn_journal_seq(&turn)?;
let belongs = imported_turns.map_or_else(
|| journal_seq > epoch.journal_seq && journal_seq < end_seq,
|turn_ids| turn_ids.contains(&turn.id),
);
belongs.then(|| WaveChatMessage {
epoch_id: epoch.id.clone(),
source: ChatMessageSource::Local { journal_seq },
turn,
})
})
.collect::<Vec<_>>();
tail_chat_messages(turns, limit)
}
pub fn committed_local_message(&self, turn: ChatTurn) -> WaveChatMessage {
let epoch = self.active_conversation_epoch();
debug_assert!(matches!(epoch.backing, ChatBacking::Local));
let journal_seq = turn_journal_seq(&turn)
.expect("a runtime-committed ChatTurn id carries its journal sequence");
WaveChatMessage {
epoch_id: epoch.id,
source: ChatMessageSource::Local { journal_seq },
turn,
}
}
pub fn discord_snapshot(&self) -> DiscordSnapshot {
let inner = self.inner();
let mut deliveries = inner
.discord_deliveries
.values()
.cloned()
.collect::<Vec<_>>();
deliveries.sort_by_key(|delivery| {
delivery
.turn_id
.strip_prefix("turn-")
.and_then(|value| value.parse::<u64>().ok())
.unwrap_or(u64::MAX)
});
DiscordSnapshot {
attachment: inner.discord.clone(),
deliveries,
}
}
pub fn try_attach_discord(
&self,
binding: DiscordChatBinding,
bot_user_id: String,
cursor: Option<String>,
) -> Result<(), JournalAppendError> {
let mut inner = self.inner();
let cursor = match inner.discord.as_ref() {
Some(attached)
if attached.binding == binding && attached.bot_user_id == bot_user_id =>
{
return Ok(())
}
Some(attached) if attached.binding == binding => attached.cursor.clone(),
_ => cursor,
};
inner
.journal
.try_append(|_| EventKind::DiscordChatAttached {
binding: binding.clone(),
bot_user_id: bot_user_id.clone(),
cursor: cursor.clone(),
})?;
inner.discord = Some(DiscordAttachment {
binding,
bot_user_id,
cursor,
});
Ok(())
}
pub fn try_deliver_discord(
&self,
text: String,
source: DiscordMessageSource,
) -> Result<bool, JournalAppendError> {
self.try_deliver_discord_input(text, source, DiscordInput::Provider)
}
pub(crate) fn try_deliver_discord_authored(
&self,
text: String,
source: DiscordMessageSource,
op: MessageOp,
) -> Result<bool, JournalAppendError> {
self.try_deliver_discord_input(text, source, DiscordInput::Authored(op))
}
fn try_deliver_discord_input(
&self,
text: String,
source: DiscordMessageSource,
input: DiscordInput,
) -> Result<bool, JournalAppendError> {
let mut inner = self.inner();
let active_binding = inner
.conversation_epochs
.last()
.and_then(|epoch| epoch.backing.discord_binding());
if active_binding.as_ref() != Some(&source.binding) {
return Ok(false);
}
if inner.messages.values().any(|known| {
known.source.as_ref().is_some_and(|known| {
known.binding == source.binding && known.message_id == source.message_id
})
}) {
return Ok(false);
}
let event = inner.journal.try_append(|seq| {
let id = MessageId(format!("msg-{seq}"));
match input {
DiscordInput::Provider => EventKind::DiscordUserMessage {
id,
text: text.clone(),
source: source.clone(),
},
DiscordInput::Authored(op) => EventKind::DiscordAuthoredMessage {
id,
op,
text: text.clone(),
source: source.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();
self.commit_locked(&mut inner, turn);
let pending = PendingMessage {
id,
op: input.op(),
text: format!("[{}]\n{}", source.uri(), text),
source: Some(source.clone()),
};
inner.messages.insert(pending.id.clone(), pending.clone());
inner.pending_messages.push(pending.clone());
let _ = self.inbox_tx.send(InboxItem::Message(pending));
Ok(true)
}
pub fn try_advance_discord_cursor(
&self,
binding: &DiscordChatBinding,
message_id: String,
) -> Result<(), JournalAppendError> {
let mut inner = self.inner();
let Some(attached) = inner.discord.as_ref() else {
return Ok(());
};
if &attached.binding != binding || attached.cursor.as_deref() == Some(&message_id) {
return Ok(());
}
inner
.journal
.try_append(|_| EventKind::DiscordChatCursorAdvanced {
binding: binding.clone(),
message_id: message_id.clone(),
})?;
if let Some(attached) = inner.discord.as_mut() {
attached.cursor = Some(message_id);
}
Ok(())
}
pub fn try_confirm_discord_part(
&self,
delivery_id: &str,
part_id: &str,
provider_message_id: String,
) -> Result<(), JournalAppendError> {
let mut inner = self.inner();
let already_confirmed = inner
.discord_deliveries
.get(delivery_id)
.and_then(|delivery| delivery.confirmed.get(part_id))
.is_some();
if already_confirmed {
return Ok(());
}
inner
.journal
.try_append(|_| EventKind::DiscordChatSendConfirmed {
delivery_id: delivery_id.to_string(),
part_id: part_id.to_string(),
provider_message_id: provider_message_id.clone(),
})?;
if let Some(delivery) = inner.discord_deliveries.get_mut(delivery_id) {
delivery
.confirmed
.insert(part_id.to_string(), provider_message_id);
}
Ok(())
}
pub fn paused(&self) -> bool {
read_wave_config(&self.repo_root, &self.name)
.and_then(|config| config.paused)
.unwrap_or(false)
}
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 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 {
epoch: inner
.conversation_epochs
.last()
.cloned()
.expect("an open runtime always has an active conversation epoch"),
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(),
pending: inner.pending_messages.clone(),
tasks: inner.tasks.clone(),
projects: inner.projects.clone(),
promotions: inner.promotions.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> {
self.try_deliver(op, text)
.expect("journal truth must accept a message before runtime delivery")
}
pub fn try_deliver(
&self,
op: MessageOp,
text: String,
) -> Result<Option<ChatTurn>, JournalAppendError> {
if op == MessageOp::Interrupt && text.trim().is_empty() {
self.deliver_interrupt();
return Ok(None);
}
self.try_deliver_message(text, op).map(Some)
}
pub fn try_deliver_authored(
&self,
op: MessageOp,
text: String,
) -> Result<Option<ChatTurn>, ChatWriteError> {
if op == MessageOp::Interrupt && text.trim().is_empty() {
self.deliver_interrupt();
return Ok(None);
}
if matches!(
self.active_conversation_epoch().backing,
ChatBacking::Discord { .. }
) {
return Err(ChatWriteError::OpenDiscord);
}
self.try_deliver(op, text).map_err(ChatWriteError::from)
}
fn try_deliver_message(
&self,
text: String,
op: MessageOp,
) -> Result<ChatTurn, JournalAppendError> {
let mut inner = self.inner();
let event = inner.journal.try_append(|seq| EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op,
text: text.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();
let turn = self.commit_locked(&mut inner, turn);
let pending = PendingMessage {
id,
op,
text,
source: None,
};
inner.messages.insert(pending.id.clone(), pending.clone());
inner.pending_messages.push(pending.clone());
let _ = self.inbox_tx.send(InboxItem::Message(pending));
Ok(turn)
}
#[cfg(test)]
pub(crate) fn fail_next_journal_append(&self, failure: JournalAppendStage) {
self.inner().journal.fail_next_append(failure);
}
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(),
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(),
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(crate) fn deliver_promotion_wake(&self, wake: PromotionWake) -> bool {
let mut inner = self.inner();
let pending = promotion_wake_message(&wake);
if inner.promotions.contains_key(&pending.id) {
return false;
}
inner.journal.append(|_| EventKind::PromotionObserved {
parent_wave_id: wake.parent_wave_id.clone(),
parent: wake.parent.clone(),
});
inner.messages.insert(pending.id.clone(), pending.clone());
inner.promotions.insert(pending.id.clone(), wake.clone());
inner.pending_messages.push(pending);
let _ = self.inbox_tx.send(InboxItem::Promotion {
parent_wave_id: wake.parent_wave_id,
parent: wake.parent,
});
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(),
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");
let committed = self.commit_locked(&mut inner, turn);
if status == Lifecycle::Completed {
self.plan_discord_delivery_locked(&mut inner, &committed, &claims);
}
}
fn plan_discord_delivery_locked(
&self,
inner: &mut Inner,
turn: &ChatTurn,
claims: &[MessageId],
) {
let Some(binding) = inner
.conversation_epochs
.last()
.and_then(|epoch| epoch.backing.discord_binding())
else {
return;
};
let Some(delivery) = build_discord_delivery(turn, claims, &inner.messages, &binding) else {
return;
};
if inner.discord_deliveries.contains_key(&delivery.delivery_id) {
return;
}
inner.journal.append(|_| EventKind::DiscordChatSendPlanned {
delivery_id: delivery.delivery_id.clone(),
turn_id: delivery.turn_id.clone(),
binding: Some(delivery.binding.clone()),
sources: delivery.sources.clone(),
parts: delivery.parts.clone(),
});
inner
.discord_deliveries
.insert(delivery.delivery_id.clone(), delivery);
}
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 initialize_conversation_epoch(
journal: &mut Journal,
epochs: &mut Vec<ConversationEpoch>,
epoch_turns: &mut HashMap<String, Vec<String>>,
turns: &[ChatTurn],
discord_turn_bindings: &HashMap<String, DiscordChatBinding>,
backing: ChatBacking,
) {
let migrating_legacy = epochs.is_empty() && !turns.is_empty();
if migrating_legacy {
let imported = legacy_conversation_epochs(turns, discord_turn_bindings);
journal.append(|_| EventKind::ConversationEpochsImported {
epochs: imported.clone(),
});
for item in imported {
epoch_turns.insert(item.epoch.id.clone(), item.turn_ids);
epochs.push(item.epoch);
}
}
if !migrating_legacy
&& epochs.last().is_some_and(|epoch| {
epoch.backing == backing && !epoch.id.starts_with("chat-epoch-legacy-")
})
{
return;
}
let number = epochs.last().map_or(1, |epoch| epoch.number + 1);
let event = journal.append(|seq| EventKind::ConversationEpochStarted {
epoch_id: format!("chat-epoch-{seq}"),
number,
backing: backing.clone(),
});
let at = event.at_rfc3339();
if let Some(previous) = epochs.last_mut() {
previous.ended_at = Some(at.clone());
}
epochs.push(ConversationEpoch {
id: format!("chat-epoch-{}", event.seq),
number,
backing,
journal_seq: event.seq,
started_at: at,
ended_at: None,
});
}
fn legacy_conversation_epochs(
turns: &[ChatTurn],
discord_turn_bindings: &HashMap<String, DiscordChatBinding>,
) -> Vec<ConversationEpochImport> {
let mut imported: Vec<ConversationEpochImport> = Vec::new();
for turn in turns {
let backing = discord_turn_bindings
.get(&turn.id)
.map(ChatBacking::discord)
.unwrap_or(ChatBacking::Local);
if let Some(active) = imported
.last_mut()
.filter(|active| active.epoch.backing == backing)
{
active.turn_ids.push(turn.id.clone());
continue;
}
let number = imported.len() as u64 + 1;
if let Some(previous) = imported.last_mut() {
previous.epoch.ended_at = Some(turn.created_at.clone());
}
imported.push(ConversationEpochImport {
epoch: ConversationEpoch {
id: format!("chat-epoch-legacy-{number}"),
number,
backing,
journal_seq: turn_journal_seq(turn).unwrap_or(1).saturating_sub(1),
started_at: turn.created_at.clone(),
ended_at: None,
},
turn_ids: vec![turn.id.clone()],
});
}
imported
}
fn turn_journal_seq(turn: &ChatTurn) -> Option<u64> {
turn.id.strip_prefix("turn-")?.parse().ok()
}
fn tail_chat_messages(
messages: Vec<WaveChatMessage>,
limit: Option<usize>,
) -> Vec<WaveChatMessage> {
let take = limit.unwrap_or(messages.len()).min(messages.len());
messages[messages.len() - take..].to_vec()
}
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 build_discord_delivery(
turn: &ChatTurn,
claims: &[MessageId],
messages: &HashMap<MessageId, PendingMessage>,
binding: &DiscordChatBinding,
) -> Option<DiscordDelivery> {
if turn.text.trim().is_empty() {
return None;
}
let sources = claims
.iter()
.filter_map(|claim| {
messages
.get(claim)
.and_then(|message| message.source.clone())
})
.filter(|source| &source.binding == binding)
.collect::<Vec<_>>();
let mut hasher = Sha256::new();
hasher.update(binding.guild_id.as_bytes());
hasher.update(binding.channel_id.as_bytes());
hasher.update(turn.id.as_bytes());
let digest = format!("{:x}", hasher.finalize());
let delivery_id = format!("discord-{}", &digest[..24]);
let parts = split_discord_content(&turn.text)
.into_iter()
.enumerate()
.map(|(index, content)| DiscordMessagePart {
part_id: format!("part-{}", index + 1),
nonce: format!("lf-{}-{index}", &digest[..16]),
content,
})
.collect();
Some(DiscordDelivery {
delivery_id,
turn_id: turn.id.clone(),
binding: binding.clone(),
sources,
parts,
confirmed: HashMap::new(),
})
}
fn split_discord_content(content: &str) -> Vec<String> {
let mut parts = Vec::new();
let mut start = 0;
let mut utf16_units = 0;
for (index, character) in content.char_indices() {
let units = character.len_utf16();
if utf16_units + units > 2_000 {
parts.push(content[start..index].to_string());
start = index;
utf16_units = 0;
}
utf16_units += units;
}
if start < content.len() {
parts.push(content[start..].to_string());
}
parts
}
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(),
body: None,
activity: None,
}
}
fn open_runtime(repo: &Path) -> Arc<WaveRuntime> {
WaveRuntime::open("ship".into(), repo.to_path_buf()).expect("open runtime")
}
fn open_discord_runtime(repo: &Path, binding: &DiscordChatBinding) -> Arc<WaveRuntime> {
WaveRuntime::open_with_backing(
"ship".into(),
repo.to_path_buf(),
ChatBacking::discord(binding),
)
.expect("open Discord 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 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 {
task_id: crate::task::TaskId::from_raw("task_example"),
issue_identifier: "INF-123".to_string(),
event_id: 7,
event: crate::task::TaskEventKind::Progress {
summary: "Task needs parent attention".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::ProjectObservation {
project_id: crate::project::ProjectId::from_raw("proj_example"),
project: "developer-efficiency".to_string(),
event_id: 8,
event: crate::project::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 consumed_promotion_stays_consumed_and_deduplicated_after_reopen() {
let tmp = tempfile::tempdir().expect("tempdir");
let wake = crate::wave::PromotionWake {
parent_wave_id: crate::id::WaveId::new(),
parent: "platform".to_string(),
};
let id = wake.inbox_id();
let rt = open_runtime(tmp.path());
assert!(rt.deliver_promotion_wake(wake.clone()));
rt.apply_resident_delta(d_opened(&[&id]));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
assert!(rt.pending_messages().is_empty());
drop(rt);
let reopened = open_runtime(tmp.path());
assert!(reopened.pending_messages().is_empty());
assert!(
!reopened.deliver_promotion_wake(wake),
"replay keeps the deterministic promotion id deduplicated"
);
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("read journal");
assert_eq!(
events
.iter()
.filter(|event| matches!(&event.kind, EventKind::PromotionObserved { .. }))
.count(),
1
);
assert!(!events
.iter()
.any(|event| matches!(&event.kind, EventKind::UserMessage { .. })));
}
#[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 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(MessageOp::Message, "to a".into());
let wave = rt.thread_snapshot();
assert_eq!(wave.len(), 2);
assert_eq!(wave[0].text, "to the wave");
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);
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"
);
}
#[test]
fn discord_chat_input_is_durable_before_cursor_and_deduplicates_after_restart() {
let tmp = tempfile::tempdir().expect("tempdir");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let source = DiscordMessageSource {
binding: binding.clone(),
message_id: "101".into(),
author_id: "human".into(),
};
let rt = open_discord_runtime(tmp.path(), &binding);
rt.try_attach_discord(binding.clone(), "bot".into(), Some("100".into()))
.expect("attach at current head");
assert!(rt
.try_deliver_discord("hello".into(), source.clone())
.expect("journal input"));
assert!(!rt
.try_deliver_discord("hello".into(), source.clone())
.expect("duplicate input"));
assert_eq!(
rt.discord_snapshot()
.attachment
.expect("attached")
.cursor
.as_deref(),
Some("100"),
"input commit does not advance the fetch cursor"
);
drop(rt);
let reopened = open_discord_runtime(tmp.path(), &binding);
assert!(!reopened
.try_deliver_discord("hello".into(), source)
.expect("refetched input"));
assert_eq!(reopened.pending_messages().len(), 1);
assert!(reopened.pending_messages()[0]
.text
.starts_with("[discord://guild/channel/101]"));
reopened
.try_advance_discord_cursor(&binding, "101".into())
.expect("commit cursor");
drop(reopened);
assert_eq!(
open_discord_runtime(tmp.path(), &binding)
.discord_snapshot()
.attachment
.expect("attached")
.cursor
.as_deref(),
Some("101")
);
}
#[test]
fn discord_app_steer_keeps_its_operation_after_restart() {
let tmp = tempfile::tempdir().expect("tempdir");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let source = DiscordMessageSource {
binding: binding.clone(),
message_id: "102".into(),
author_id: "bot".into(),
};
let rt = open_discord_runtime(tmp.path(), &binding);
rt.try_attach_discord(binding.clone(), "bot".into(), Some("101".into()))
.expect("attach at current head");
assert!(rt
.try_deliver_discord_authored("change course".into(), source.clone(), MessageOp::Steer)
.expect("journal app steer"));
assert!(!rt
.try_deliver_discord_authored("change course".into(), source.clone(), MessageOp::Steer)
.expect("deduplicate app steer"));
assert_eq!(rt.pending_messages()[0].op, MessageOp::Steer);
drop(rt);
let reopened = open_discord_runtime(tmp.path(), &binding);
assert_eq!(reopened.pending_messages().len(), 1);
assert_eq!(reopened.pending_messages()[0].op, MessageOp::Steer);
assert_eq!(
reopened.pending_messages()[0]
.source
.as_ref()
.expect("Discord source"),
&source
);
}
#[test]
fn discord_chat_answer_is_planned_in_chunks_before_receipts() {
let tmp = tempfile::tempdir().expect("tempdir");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let rt = open_discord_runtime(tmp.path(), &binding);
rt.try_attach_discord(binding.clone(), "bot".into(), None)
.expect("attach");
rt.try_deliver_discord(
"question".into(),
DiscordMessageSource {
binding: binding.clone(),
message_id: "101".into(),
author_id: "human".into(),
},
)
.expect("deliver");
let message_id = rt.pending_messages()[0].id.0.clone();
rt.apply_resident_delta(d_opened(&[&message_id]));
rt.apply_resident_delta(d_text(&"x".repeat(2_001)));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let delivery = rt
.discord_snapshot()
.deliveries
.into_iter()
.next()
.expect("send intent");
assert_eq!(delivery.parts.len(), 2);
assert_eq!(delivery.parts[0].content.chars().count(), 2_000);
assert!(delivery.parts.iter().all(|part| part.nonce.len() <= 25));
assert!(delivery.confirmed.is_empty());
rt.try_confirm_discord_part(
&delivery.delivery_id,
&delivery.parts[0].part_id,
"provider-1".into(),
)
.expect("confirm first part");
drop(rt);
let reopened = open_discord_runtime(tmp.path(), &binding);
let resumed = &reopened.discord_snapshot().deliveries[0];
assert_eq!(resumed.confirmed.len(), 1);
assert_eq!(resumed.parts.len(), 2);
}
#[test]
fn discord_epoch_routes_autonomous_agent_speech_as_a_top_level_message() {
let tmp = tempfile::tempdir().expect("tempdir");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let rt = open_discord_runtime(tmp.path(), &binding);
rt.try_attach_discord(binding.clone(), "bot".into(), None)
.expect("attach");
rt.apply_resident_delta(d_opened(&[]));
rt.apply_resident_delta(d_text("autonomous update"));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let delivery = rt
.discord_snapshot()
.deliveries
.into_iter()
.next()
.expect("active Discord epoch chooses delivery");
assert_eq!(delivery.binding, binding);
assert!(
delivery.sources.is_empty(),
"no reply target means top-level"
);
assert_eq!(delivery.parts[0].content, "autonomous update");
}
#[test]
fn discord_epoch_delivers_speech_that_claims_a_typed_task_observation() {
let tmp = tempfile::tempdir().expect("tempdir");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let rt = open_discord_runtime(tmp.path(), &binding);
let observation = crate::task::TaskObservation {
task_id: crate::task::TaskId::from_raw("task_example"),
issue_identifier: "INF-123".into(),
event_id: 7,
event: crate::task::TaskEventKind::Progress {
summary: "Task needs parent attention".into(),
},
};
assert!(rt.deliver_task_observation(observation.clone()));
rt.apply_resident_delta(d_opened(&[&observation.inbox_id()]));
rt.apply_resident_delta(d_text("I handled the child update."));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
let delivery = rt
.discord_snapshot()
.deliveries
.into_iter()
.next()
.expect("typed input cannot suppress active-backing speech");
assert_eq!(delivery.binding, binding);
assert!(
delivery.sources.is_empty(),
"typed input is not a reply target"
);
assert_eq!(delivery.parts[0].content, "I handled the child update.");
}
#[test]
fn wave_chat_backing_switch_rejects_parallel_local_compose() {
let tmp = tempfile::tempdir().expect("tempdir");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let local = open_runtime(tmp.path());
let local_epoch = local.active_conversation_epoch();
assert_eq!(local_epoch.number, 1);
assert_eq!(local_epoch.backing, ChatBacking::Local);
local
.try_deliver_authored(MessageOp::Message, "local question".into())
.expect("local write")
.expect("local turn");
let local_messages = local.chat_messages(None, None);
assert_eq!(local_messages.len(), 1);
assert!(matches!(
local_messages[0].source,
ChatMessageSource::Local { .. }
));
drop(local);
let discord = open_discord_runtime(tmp.path(), &binding);
let discord_epoch = discord.active_conversation_epoch();
assert_eq!(discord_epoch.number, 2);
assert_eq!(discord_epoch.backing, ChatBacking::discord(&binding));
let epochs = discord.conversation_epochs();
assert_eq!(epochs.len(), 2);
assert!(epochs[0].ended_at.is_some());
assert_eq!(
discord.chat_messages(Some(&local_epoch.id), None),
local_messages,
"the earlier local epoch remains byte-identical"
);
let before = crate::wave::journal::read_events(&journal_path(tmp.path(), "ship"));
let before_pending = discord.pending_messages();
let error = discord
.try_deliver_authored(MessageOp::Message, "shadow message".into())
.expect_err("Discord mode rejects Loopflow compose");
assert!(matches!(error, ChatWriteError::OpenDiscord));
let after = crate::wave::journal::read_events(&journal_path(tmp.path(), "ship"));
assert_eq!(before, after, "rejection appends no local journal event");
assert_eq!(before_pending, discord.pending_messages());
assert!(discord.chat_messages(None, None).is_empty());
discord
.try_deliver_authored(MessageOp::Interrupt, String::new())
.expect("bare interrupt remains available");
drop(discord);
let reopened = open_discord_runtime(tmp.path(), &binding);
assert_eq!(
reopened.active_conversation_epoch().id,
discord_epoch.id,
"a restart with the same backing resumes the epoch"
);
drop(reopened);
let local_again = open_runtime(tmp.path());
assert_eq!(local_again.active_conversation_epoch().number, 3);
assert_eq!(
local_again.active_conversation_epoch().backing,
ChatBacking::Local
);
assert_eq!(local_again.conversation_epochs().len(), 3);
}
#[test]
fn legacy_local_epoch_survives_the_migration_restart() {
let tmp = tempfile::tempdir().expect("tempdir");
let path = journal_path(tmp.path(), "ship");
let (mut journal, _) = Journal::open(&path).expect("legacy journal");
journal.append(|_| EventKind::UserMessage {
id: MessageId("legacy-message".into()),
op: MessageOp::Message,
text: "before epochs".into(),
});
drop(journal);
let migrated = open_runtime(tmp.path());
let epochs = migrated.conversation_epochs();
assert_eq!(epochs.len(), 2);
assert_eq!(epochs[0].id, "chat-epoch-legacy-1");
assert_eq!(epochs[0].backing, ChatBacking::Local);
let legacy_messages = migrated.chat_messages(Some(&epochs[0].id), None);
assert_eq!(legacy_messages.len(), 1);
assert_eq!(legacy_messages[0].turn.text, "before epochs");
drop(migrated);
let reopened = open_runtime(tmp.path());
assert_eq!(reopened.conversation_epochs(), epochs);
assert_eq!(
reopened.chat_messages(Some("chat-epoch-legacy-1"), None),
legacy_messages
);
}
#[test]
fn legacy_mixed_chat_imports_truthful_backing_epochs_atomically() {
let tmp = tempfile::tempdir().expect("tempdir");
let path = journal_path(tmp.path(), "ship");
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let (mut journal, _) = Journal::open(&path).expect("legacy journal");
journal.append(|_| EventKind::UserMessage {
id: MessageId("local-message".into()),
op: MessageOp::Message,
text: "local history".into(),
});
journal.append(|_| EventKind::DiscordUserMessage {
id: MessageId("discord-message".into()),
text: "provider history".into(),
source: DiscordMessageSource {
binding: binding.clone(),
message_id: "101".into(),
author_id: "human".into(),
},
});
drop(journal);
let migrated = open_runtime(tmp.path());
let epochs = migrated.conversation_epochs();
assert_eq!(epochs.len(), 3);
assert_eq!(epochs[0].backing, ChatBacking::Local);
assert_eq!(epochs[1].backing, ChatBacking::discord(&binding));
assert_eq!(epochs[2].backing, ChatBacking::Local);
assert_eq!(
migrated.chat_messages(Some(&epochs[0].id), None)[0]
.turn
.text,
"local history"
);
assert!(
migrated.chat_messages(Some(&epochs[1].id), None).is_empty(),
"Discord history remains provider-projected, never mislabeled local"
);
let imported = crate::wave::journal::read_events(&path)
.into_iter()
.filter(|event| matches!(event.kind, EventKind::ConversationEpochsImported { .. }))
.count();
assert_eq!(imported, 1, "the entire legacy catalog is one append");
drop(migrated);
assert_eq!(open_runtime(tmp.path()).conversation_epochs(), epochs);
}
}