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::BodyProvenance;
use crate::chat::turns::{ChatRole, ChatTurn, TurnDelta};
use crate::chat::types::{ConversationItem, Lifecycle};
use crate::controller::wave::channel::{Author, Message};
use crate::controller::wave::chat::{
ChatBacking, ChatMessageSource, ConversationEpoch, WaveChatMessage,
};
#[cfg(test)]
use crate::controller::wave::journal::JournalAppendStage;
use crate::controller::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,
};
use crate::controller::wave::state::{can_transition, LoopState};
use crate::controller::wave::wire::{ProviderSessionRef, ResidentDelta, ResidentStateTo};
use crate::work::project::ProjectObservation;
use crate::work::task::TaskObservation;
use crate::work::wave::PromotionWake;
const TURN_BROADCAST_CAPACITY: usize = 256;
const STATE_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,
}
#[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 pending: Vec<PendingMessage>,
pub chat_tail: 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,
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>,
state: LoopState,
pending_messages: Vec<PendingMessage>,
messages: HashMap<MessageId, PendingMessage>,
discord: Option<DiscordAttachment>,
discord_deliveries: HashMap<String, DiscordDelivery>,
chat_reply_targets: HashMap<String, MessageId>,
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>,
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, mut events) = Journal::open(&journal_path(&repo_root, &name))?;
crate::controller::wave::recovery::prepare(&mut journal, &mut events, &name)?;
let mut fold = fold_thread(&events);
if let Some(body) = fold.open.iter().find_map(|turn| turn.body.as_ref()) {
anyhow::bail!("Wave '{name}' has an unfinished attempt {} with unknown provider liveness; preserve {} until its owner records termination", body.body_id, journal_path(&repo_root, &name).display());
}
const ABANDONED: &str = "startup janitor: turn abandoned by server restart";
for mut turn in fold.open {
let finished = journal.append(|_| EventKind::TurnFinished {
turn_id: turn.id.clone(),
status: Lifecycle::Failed,
termination_reason: Some(ABANDONED.to_string()),
});
turn.status = Lifecycle::Failed;
turn.close_body(finished.at_rfc3339(), Some(ABANDONED.to_string()));
fold.turns.push(turn);
}
restore_pending(
&mut fold.pending_messages,
&fold.messages,
&fold.open_claims,
);
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");
let completed_turns = fold
.turns
.iter()
.filter(|turn| turn.role == ChatRole::Assistant && turn.status == Lifecycle::Completed)
.filter(|turn| !planned_turns.contains(&turn.id))
.filter(|turn| turn_journal_seq(turn).is_some_and(|seq| seq > active_epoch.journal_seq))
.cloned()
.collect::<Vec<_>>();
for turn in &completed_turns {
let Some(binding) = active_epoch.backing.discord_binding() else {
continue;
};
let Some(delivery) = build_discord_delivery(
turn,
fold.chat_reply_targets.get(&turn.id),
&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 (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,
state,
pending_messages: fold.pending_messages,
messages: fold.messages,
discord: fold.discord,
discord_deliveries: fold.discord_deliveries,
chat_reply_targets: fold.chat_reply_targets,
tasks: fold.tasks,
projects: fold.projects,
promotions: fold.promotions,
}),
turn_tx,
state_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.author_name = source.author_name.clone();
turn.created_at = event.at_rfc3339();
self.commit_locked(&mut inner, turn);
let pending = PendingMessage {
id,
op: input.op(),
text: super::journal::attributed_message(
&format!("[{}]\n{}", source.uri(), text),
source.author_name.as_deref(),
),
source: Some(source.clone()),
};
inner.messages.insert(pending.id.clone(), 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(())
}
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 read_channel(&self, since: Option<u64>) -> Vec<Message> {
let inner = self.inner();
snapshot_tail_locked(&inner, None)
.into_iter()
.filter_map(Message::from_turn)
.filter(|(seq, _)| since.is_none_or(|cursor| *seq > cursor))
.map(|(seq, mut message)| {
let source = inner
.messages
.get(&MessageId(format!("msg-{seq}")))
.and_then(|pending| pending.source.as_ref());
if let Some(source) = source {
message.author = Author::Bridge {
platform: "discord".to_string(),
user: source.author_id.clone(),
name: source.author_name.clone(),
};
}
if message.is_own() {
message.reply_to = inner
.chat_reply_targets
.get(&message.id)
.and_then(channel_turn_id_for_message);
}
message
})
.collect()
}
pub fn unanswered_chat_tail(&self) -> Vec<PendingMessage> {
unanswered_chat_tail_locked(&self.inner())
}
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()
}
fn update_body_session(&self, body_id: &str, session_id: &str) {
let mut inner = self.inner();
let Some(open) = inner.open.as_mut() else {
return;
};
let Some(body) = open
.turn
.body
.as_mut()
.filter(|body| body.body_id == body_id)
else {
return;
};
body.session_id = Some(session_id.to_string());
let turn = open.turn.clone();
inner.journal.append(|_| EventKind::BodySessionUpdated {
body_id: body_id.to_string(),
session_id: session_id.to_string(),
});
let _ = self
.turn_tx
.send(TurnBroadcast::Whole(TurnFrame::share(turn)));
}
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(),
pending: inner.pending_messages.clone(),
chat_tail: unanswered_chat_tail_locked(&inner),
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();
if inner.open.is_none() {
return false;
}
self.finish_turn_locked(&mut inner, status, Some(reason.to_string()));
true
}
fn commit_locked(&self, inner: &mut Inner, turn: ChatTurn) -> ChatTurn {
turn.validate()
.expect("Wave thread entries must satisfy the ChatTurn wire invariant");
inner.thread.push(turn.clone());
let _ = self
.turn_tx
.send(TurnBroadcast::Whole(TurnFrame::share(turn.clone())));
turn
}
pub fn try_deliver(
&self,
op: MessageOp,
text: String,
author_name: Option<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);
}
let mut inner = self.inner();
let author_name = author_name
.as_deref()
.and_then(crate::engine::config::normalize_user_name);
let event = inner.journal.try_append(|seq| EventKind::UserMessage {
id: MessageId(format!("msg-{seq}")),
op,
text: text.clone(),
author_name: author_name.clone(),
})?;
let id = MessageId(format!("msg-{}", event.seq));
let mut turn = ChatTurn::user(format!("turn-{}", event.seq), text.clone());
turn.author_name = author_name.clone();
turn.created_at = event.at_rfc3339();
let turn = self.commit_locked(&mut inner, turn);
let pending = PendingMessage {
id,
op,
text: super::journal::attributed_message(&text, author_name.as_deref()),
source: None,
};
inner.messages.insert(pending.id.clone(), pending.clone());
let _ = self.inbox_tx.send(InboxItem::Message(pending));
Ok(Some(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_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>,
reply_to: Option<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(),
});
}
if let Some(message_id) = &reply_to {
inner.journal.append(|_| EventKind::ChatReplyLinked {
turn_id: turn_id.clone(),
message_id: message_id.clone(),
});
}
inner.journal.append(|_| EventKind::TurnFinished {
turn_id: turn_id.clone(),
status: turn.status,
termination_reason: None,
});
let committed = ChatTurn {
id: turn_id,
created_at: started.at_rfc3339(),
..turn
};
if let Some(message_id) = reply_to {
inner
.chat_reply_targets
.insert(committed.id.clone(), message_id);
}
let committed = self.commit_locked(&mut inner, committed);
self.plan_discord_delivery_locked(&mut inner, &committed);
committed
}
pub fn apply_resident_delta(&self, delta: ResidentDelta) {
match delta {
ResidentDelta::TurnOpened { answers, body } => self.resident_turn_opened(answers, body),
ResidentDelta::TurnText { text } => self.resident_turn_text(text),
ResidentDelta::ChatReply { message_id, text } => {
self.resident_chat_reply(message_id, text)
}
ResidentDelta::TurnItem { item } => self.resident_turn_item(item),
ResidentDelta::TurnFinished { status, reason } => {
self.finish_turn_locked(&mut self.inner(), status, reason)
}
ResidentDelta::BodySessionUpdated {
body_id,
session_id,
} => {
self.update_body_session(&body_id, &session_id);
}
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>, body: Option<BodyProvenance>) {
let mut inner = self.inner();
if inner.open.is_some() {
self.finish_turn_locked(
&mut inner,
Lifecycle::Failed,
Some("stale open turn closed".into()),
);
}
let answers = claim_answers(&mut inner, answers);
let claims = answers.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 {
author_name: None,
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,
text_items: 0,
claims,
});
}
fn resident_chat_reply(&self, message_id: String, text: String) {
let message_id = MessageId(message_id);
if !self.inner().messages.contains_key(&message_id) {
return;
}
self.append_finalized_turn(
ChatTurn {
id: String::new(),
role: ChatRole::Assistant,
author_name: None,
text,
status: Lifecycle::Completed,
items: Vec::new(),
created_at: String::new(),
body: None,
activity: None,
},
Vec::new(),
Some(message_id),
);
}
fn resident_turn_text(&self, text: String) {
let mut inner = self.inner();
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.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 finish_turn_locked(&self, inner: &mut Inner, status: Lifecycle, reason: Option<String>) {
let Some(OpenTurn {
mut turn, claims, ..
}) = inner.open.take()
else {
tracing::warn!("TurnFinished with no open turn; dropped");
return;
};
let finished = inner.journal.append(|_| EventKind::TurnFinished {
turn_id: turn.id.clone(),
status,
termination_reason: reason.clone(),
});
if status != Lifecycle::Completed {
restore_pending(&mut inner.pending_messages, &inner.messages, &claims);
}
turn.status = status;
turn.close_body(finished.at_rfc3339(), reason);
self.transition_locked(inner, LoopState::Idle, "turn finalized");
let committed = self.commit_locked(inner, turn);
if status == Lifecycle::Completed {
self.plan_discord_delivery_locked(inner, &committed);
}
}
fn plan_discord_delivery_locked(&self, inner: &mut Inner, turn: &ChatTurn) {
let Some(binding) = inner
.conversation_epochs
.last()
.and_then(|epoch| epoch.backing.discord_binding())
else {
return;
};
let Some(delivery) = build_discord_delivery(
turn,
inner.chat_reply_targets.get(&turn.id),
&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 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 unanswered_chat_tail_locked(inner: &Inner) -> Vec<PendingMessage> {
let message_seq = |message: &PendingMessage| {
message
.id
.0
.strip_prefix("msg-")
.and_then(|seq| seq.parse::<u64>().ok())
};
inner
.messages
.values()
.filter(|message| message.op == MessageOp::Message)
.filter(|message| {
!inner
.chat_reply_targets
.values()
.any(|reply_to| reply_to == &message.id)
})
.filter(|message| message_seq(message).is_some())
.max_by_key(|message| message_seq(message).unwrap_or_default())
.cloned()
.into_iter()
.collect()
}
fn channel_turn_id_for_message(message_id: &MessageId) -> Option<String> {
message_id
.0
.strip_prefix("msg-")
.map(|seq| format!("turn-{seq}"))
}
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,
reply_to: Option<&MessageId>,
messages: &HashMap<MessageId, PendingMessage>,
binding: &DiscordChatBinding,
) -> Option<DiscordDelivery> {
if turn.text.trim().is_empty() {
return None;
}
let sources = reply_to
.and_then(|message_id| messages.get(message_id))
.and_then(|message| message.source.clone())
.filter(|source| &source.binding == binding)
.into_iter()
.collect();
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
}
#[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 {
author_name: None,
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 deliver_wake(rt: &WaveRuntime, n: i64, summary: &str) -> MessageId {
let observation = crate::work::task::TaskObservation {
task_id: crate::work::task::TaskId::from_raw(format!("task_{n}")),
issue_identifier: format!("INF-{n}"),
event_id: n,
event: crate::work::task::TaskEventKind::Progress {
summary: summary.to_string(),
},
};
assert!(rt.deliver_task_observation(observation.clone()));
MessageId(observation.inbox_id())
}
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")
}
#[test]
fn read_channel_returns_the_conversation_including_the_bots_own_posts() {
use crate::controller::wave::channel::Author;
let tmp = tempfile::tempdir().expect("tempdir");
let runtime = open_runtime(tmp.path());
runtime
.try_deliver(MessageOp::Message, "what's the top task?".into(), None)
.expect("journal write")
.expect("chat message");
runtime.append_finalized_turn(progress_turn("The top task is LOO-258."), Vec::new(), None);
let all = runtime.read_channel(None);
assert_eq!(all.len(), 2);
assert!(matches!(all[0].author, Author::Human { .. }));
assert_eq!(all[0].content, "what's the top task?");
assert!(all[1].is_own());
assert_eq!(all[1].content, "The top task is LOO-258.");
let cursor = all[0]
.id
.strip_prefix("turn-")
.and_then(|seq| seq.parse::<u64>().ok())
.expect("turn seq");
let after = runtime.read_channel(Some(cursor));
assert_eq!(after.len(), 1);
assert!(after[0].is_own());
}
#[test]
fn a_later_bot_post_never_hides_the_restart_chat_trigger() {
let tmp = tempfile::tempdir().expect("tempdir");
let runtime = open_runtime(tmp.path());
runtime
.try_deliver(MessageOp::Message, "please answer this".into(), None)
.expect("journal write")
.expect("chat message");
runtime.append_finalized_turn(progress_turn("unrelated task finished"), Vec::new(), None);
let triggers = runtime.unanswered_chat_tail();
assert_eq!(triggers.len(), 1);
assert_eq!(triggers[0].text, "please answer this");
}
#[test]
fn explicit_chat_reply_relation_survives_restart_without_consumption() {
let tmp = tempfile::tempdir().expect("tempdir");
let runtime = open_runtime(tmp.path());
runtime
.try_deliver(MessageOp::Message, "what is two plus two?".into(), None)
.expect("journal write")
.expect("chat message");
let message_id = runtime.unanswered_chat_tail()[0].id.0.clone();
runtime.apply_resident_delta(ResidentDelta::ChatReply {
message_id,
text: "Four.".into(),
});
assert!(runtime.unanswered_chat_tail().is_empty());
let channel = runtime.read_channel(None);
assert_eq!(channel[1].reply_to.as_deref(), Some(channel[0].id.as_str()));
drop(runtime);
let reopened = open_runtime(tmp.path());
assert!(reopened.unanswered_chat_tail().is_empty());
let channel = reopened.read_channel(None);
assert_eq!(channel[1].reply_to.as_deref(), Some(channel[0].id.as_str()));
}
#[test]
fn chat_reply_preserves_active_governance_and_terminal_requeues_once() {
let tmp = tempfile::tempdir().unwrap();
let runtime = open_runtime(tmp.path());
runtime
.try_deliver(MessageOp::Message, "status?".into(), None)
.unwrap();
let message = runtime.unanswered_chat_tail()[0].id.0.clone();
let wake = deliver_wake(&runtime, 1, "Task finished");
let body = BodyProvenance::for_wave(tmp.path());
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: vec![wake.0.clone()],
body: Some(body.clone()),
});
runtime.apply_resident_delta(d_text("before "));
let state = runtime.loop_state();
runtime.apply_resident_delta(ResidentDelta::ChatReply {
message_id: message.clone(),
text: "Still working.".into(),
});
assert_eq!(runtime.loop_state(), state);
assert!(runtime.subscribe_with_snapshot(None).pending.is_empty());
runtime.apply_resident_delta(ResidentDelta::BodySessionUpdated {
body_id: body.body_id.clone(),
session_id: "provider-session".into(),
});
runtime.apply_resident_delta(d_text("after"));
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Failed,
reason: Some("capacity".into()),
});
let terminal = runtime
.thread_snapshot()
.into_iter()
.find(|t| t.body.is_some())
.unwrap();
assert_eq!(terminal.text, "before after");
let provenance = terminal.body.as_ref().unwrap();
assert_eq!(provenance.session_id.as_deref(), Some("provider-session"));
assert_eq!(provenance.termination_reason.as_deref(), Some("capacity"));
assert!(provenance.ended_at.is_some());
runtime.apply_resident_delta(d_finished(Lifecycle::Failed));
assert_eq!(runtime.subscribe_with_snapshot(None).pending.len(), 1);
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).unwrap();
assert_eq!(
events
.iter()
.filter(|event| matches!(&event.kind,
EventKind::TurnFinished { turn_id, .. } if turn_id == &terminal.id))
.count(),
1
);
assert!(!events
.iter()
.any(|event| matches!(event.kind, EventKind::MessagesRequeued { .. })));
let fold = fold_thread(&events);
assert!(fold.open.is_empty());
assert_eq!(fold.pending_messages.len(), 1);
assert_eq!(fold.pending_messages[0].id, wake);
assert_eq!(
fold.turns.iter().find(|t| t.id == terminal.id),
Some(&terminal)
);
}
fn d_opened(answers: &[&str]) -> ResidentDelta {
ResidentDelta::TurnOpened {
body: None,
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_finished(status: Lifecycle) -> ResidentDelta {
ResidentDelta::TurnFinished {
status,
reason: None,
}
}
fn complete_body(rt: &WaveRuntime, harness: &str, session_id: Option<&str>) {
let mut body = BodyProvenance::for_wave(rt.repo_root());
body.harness = Some(harness.to_string());
body.session_id = session_id.map(str::to_string);
rt.apply_resident_delta(ResidentDelta::TurnOpened {
body: Some(body),
answers: vec![],
});
rt.apply_resident_delta(d_text("done"));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
}
#[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(), None);
let b = rt.append_finalized_turn(progress_turn("two"), Vec::new(), None);
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
.try_deliver(MessageOp::Message, "how goes it?".into(), None)
.expect("journal write")
.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!(rt.pending_messages().is_empty());
}
#[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::work::task::TaskObservation {
task_id: crate::work::task::TaskId::from_raw("task_example"),
issue_identifier: "INF-123".to_string(),
event_id: 7,
event: crate::work::task::TaskEventKind::Progress {
summary: "Task is waiting on its parent".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.try_deliver(MessageOp::Message, format!("message {i}"), None)
.expect("journal write")
.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.try_deliver(MessageOp::Message, "message 5".into(), None)
.expect("journal write")
.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::work::project::ProjectObservation {
project_id: crate::work::project::ProjectId::from_raw("proj_example"),
project: "developer-efficiency".to_string(),
event_id: 8,
event: crate::work::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::work::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(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
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");
assert!(events
.iter()
.any(|event| matches!(event.kind, EventKind::TurnFinished { .. })));
}
#[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 = deliver_wake(&rt, 1, "first");
let m2 = deliver_wake(&rt, 2, "second");
assert_eq!(rt.pending_messages().len(), 2);
rt.apply_resident_delta(d_opened(&[&m1.0, &m2.0, "msg-999"]));
rt.apply_resident_delta(d_finished(Lifecycle::Completed));
assert!(
rt.pending_messages().is_empty(),
"claimed wakes leave the live pending fold"
);
rt.apply_resident_delta(d_opened(&[&m1.0]));
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![m1.clone(), m2.clone()],
"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_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 mut body = BodyProvenance::for_wave(tmp.path());
body.body_id = "body-dead".into();
body.session_id = Some("session-dead".into());
body.harness = Some("codex".into());
rt.apply_resident_delta(ResidentDelta::TurnOpened {
body: Some(body),
answers: vec![],
});
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 (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let fold = crate::controller::wave::journal::fold_thread(&events);
assert!(fold.open.is_empty());
assert_eq!(fold.turns.last().unwrap().status, Lifecycle::Interrupted);
assert!(fold.playhead.is_none());
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_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 = deliver_wake(&rt, 1, "do the thing");
rt.apply_resident_delta(d_opened(&[&m1.0]));
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].id, m1);
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 = deliver_wake(&rt3, 2, "second");
rt3.apply_resident_delta(d_opened(&[&m2.0]));
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 = deliver_wake(&rt, 1, "answer me");
rt.apply_resident_delta(d_opened(&[&m.0]));
assert!(rt.pending_messages().is_empty());
m.0
};
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 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.try_deliver(MessageOp::Message, "to the wave".into(), None)
.expect("journal write")
.expect("user turn");
rt.try_deliver(MessageOp::Message, "to a".into(), None)
.expect("journal write");
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::controller::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 texts: Vec<_> = rt2
.read_channel(None)
.into_iter()
.map(|message| message.content)
.collect();
assert_eq!(texts.len(), 2);
assert!(texts.contains(&"to the wave".to_string()) && texts.contains(&"to a".to_string()));
}
#[test]
fn subscription_carries_pending_replay_and_live_inbox() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let before = deliver_wake(&rt, 1, "before");
let mut sub = rt.subscribe_with_snapshot(None);
assert_eq!(sub.pending.len(), 1);
assert_eq!(sub.pending[0].id, before);
assert!(sub.inbox_rx.try_recv().is_err(), "no frames from before");
rt.try_deliver(MessageOp::Message, "after".into(), None)
.expect("journal write")
.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.try_deliver(MessageOp::Message, format!("m-{writer}-{i}"), None)
.expect("journal write")
.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 {
author_name: None,
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"));
let channel = reopened.read_channel(None);
assert_eq!(channel.len(), 1, "one readable discord message");
assert_eq!(channel[0].content, "hello");
assert!(matches!(
&channel[0].author,
Author::Bridge { platform, user, .. }
if platform == "discord" && user == "human"
));
assert!(reopened.pending_messages().is_empty());
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 request_authors_survive_local_replay_without_rewriting_conversation() {
let tmp = tempfile::tempdir().unwrap();
let rt = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).unwrap();
for name in [Some("Jack"), Some("Maya"), None] {
rt.try_deliver(
MessageOp::Message,
"You chose the prototype path.".into(),
name.map(str::to_string),
)
.unwrap();
}
let live = rt.read_channel(None);
let turns = rt.thread_tail(None);
drop(rt);
let reopened = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).unwrap();
assert_eq!(reopened.read_channel(None), live);
assert_eq!(reopened.thread_tail(None), turns);
assert_eq!(
turns
.iter()
.map(|turn| turn.author_name.as_deref())
.collect::<Vec<_>>(),
vec![Some("Jack"), Some("Maya"), None]
);
assert!(turns
.iter()
.all(|turn| turn.text == "You chose the prototype path."));
assert_eq!(
live[0].author,
Author::Human {
name: "Jack".into()
}
);
assert_eq!(
live[1].author,
Author::Human {
name: "Maya".into()
}
);
assert_eq!(
live[2].author,
Author::Human {
name: String::new()
}
);
}
#[test]
fn request_authors_keep_discord_ids_separate_from_display_names() {
let tmp = tempfile::tempdir().unwrap();
let binding = DiscordChatBinding {
guild_id: "guild".into(),
channel_id: "channel".into(),
};
let rt = open_discord_runtime(tmp.path(), &binding);
for (id, name) in [("101", Some("Jack")), ("102", Some("Maya")), ("103", None)] {
rt.try_deliver_discord(
"prototype".into(),
DiscordMessageSource {
binding: binding.clone(),
message_id: id.into(),
author_id: format!("person-{id}"),
author_name: name.map(str::to_string),
},
)
.unwrap();
}
let live = rt.read_channel(None);
drop(rt);
let reopened = open_discord_runtime(tmp.path(), &binding);
assert_eq!(reopened.read_channel(None), live);
for (message, (id, name)) in
live.iter()
.zip([("101", Some("Jack")), ("102", Some("Maya")), ("103", None)])
{
assert_eq!(
message.author,
Author::Bridge {
platform: "discord".into(),
user: format!("person-{id}"),
name: name.map(str::to_string)
}
);
}
}
#[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 {
author_name: None,
binding: binding.clone(),
message_id: "101".into(),
author_id: "human".into(),
},
)
.expect("deliver");
rt.apply_resident_delta(d_opened(&[]));
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::work::task::TaskObservation {
task_id: crate::work::task::TaskId::from_raw("task_example"),
issue_identifier: "INF-123".into(),
event_id: 7,
event: crate::work::task::TaskEventKind::Progress {
summary: "Task is waiting on its parent".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(MessageOp::Message, "local question".into(), None)
.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::controller::wave::journal::read_events(&journal_path(tmp.path(), "ship"));
let before_pending = discord.pending_messages();
let error = discord
.try_deliver(MessageOp::Message, "shadow message".into(), None)
.expect_err("Discord mode rejects Loopflow compose");
assert!(matches!(error, ChatWriteError::OpenDiscord));
let after =
crate::controller::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(MessageOp::Interrupt, String::new(), None)
.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 {
author_name: None,
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 {
author_name: None,
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 {
author_name: None,
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::controller::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);
}
}