use std::sync::Arc;
use std::time::Duration;
use tokio::process::Child;
use tokio::sync::{broadcast, mpsc};
use tokio::time::Instant;
use crate::chat::turns::{ChatRole, ChatTurn};
use crate::chat::types::Lifecycle;
use crate::wave::journal::MessageOp;
use crate::wave::registry::process_alive;
use crate::wave::runtime::{InboxItem, TurnFrame, WaveRuntime};
use crate::wave::server::ResidentDoor;
use crate::wave::state::MindState;
pub const RESPAWN_BACKOFF: [Duration; 3] = [
Duration::from_secs(5 * 60),
Duration::from_secs(15 * 60),
Duration::from_secs(45 * 60),
];
pub const LISTENER_INTERRUPT_DEADLINE: Duration = Duration::from_secs(20);
pub const ATTACH_PROBE: Duration = Duration::from_secs(10);
pub type SpawnResident = Box<dyn FnMut() -> std::io::Result<Child> + Send>;
#[derive(Debug, Clone)]
pub struct SupervisorConfig {
pub respawn_backoff: Vec<Duration>,
pub interrupt_deadline: Duration,
pub attach_probe: Duration,
}
impl Default for SupervisorConfig {
fn default() -> Self {
Self {
respawn_backoff: RESPAWN_BACKOFF.to_vec(),
interrupt_deadline: LISTENER_INTERRUPT_DEADLINE,
attach_probe: ATTACH_PROBE,
}
}
}
fn respawn_delay(backoff: &[Duration], attempts_made: u32) -> Option<Duration> {
backoff
.get(attempts_made as usize)
.or(backoff.last())
.copied()
}
async fn wait_child(child: &mut Option<Child>) -> std::io::Result<std::process::ExitStatus> {
child
.as_mut()
.expect("guarded by if in select")
.wait()
.await
}
pub(crate) async fn sleep_until_opt(deadline: Option<Instant>) {
match deadline {
Some(deadline) => tokio::time::sleep_until(deadline).await,
None => std::future::pending().await,
}
}
#[derive(Debug, Clone)]
pub struct SupervisorHandle {
attach_tx: mpsc::UnboundedSender<u32>,
}
impl SupervisorHandle {
pub fn on_attach(&self, pid: u32) {
let _ = self.attach_tx.send(pid);
}
}
pub struct Supervisor {
runtime: Arc<WaveRuntime>,
door: ResidentDoor,
spawner: Option<SpawnResident>,
config: SupervisorConfig,
inbox_rx: broadcast::Receiver<InboxItem>,
state_rx: broadcast::Receiver<MindState>,
turn_rx: broadcast::Receiver<Arc<TurnFrame>>,
attach_rx: mpsc::UnboundedReceiver<u32>,
attach_tx: mpsc::UnboundedSender<u32>,
child: Option<Child>,
respawn_at: Option<Instant>,
attempts: u32,
interrupt_at: Option<Instant>,
}
impl Supervisor {
pub fn new(
runtime: Arc<WaveRuntime>,
door: ResidentDoor,
spawner: Option<SpawnResident>,
config: SupervisorConfig,
) -> Self {
let (attach_tx, attach_rx) = mpsc::unbounded_channel();
Self {
inbox_rx: runtime.subscribe_inbox(),
state_rx: runtime.subscribe_states(),
turn_rx: runtime.subscribe_turns(),
runtime,
door,
spawner,
config,
attach_rx,
attach_tx,
child: None,
respawn_at: None,
attempts: 0,
interrupt_at: None,
}
}
pub fn handle(&self) -> SupervisorHandle {
SupervisorHandle {
attach_tx: self.attach_tx.clone(),
}
}
pub async fn run(mut self) {
if self.spawner.is_some() {
self.spawn().await;
}
let mut probe = tokio::time::interval(self.config.attach_probe);
probe.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
tokio::select! {
status = wait_child(&mut self.child), if self.child.is_some() => {
let reason = match status {
Ok(status) => format!("resident process exited ({status})"),
Err(err) => format!("resident process unwaitable: {err}"),
};
self.child = None;
self.on_resident_death(&reason);
}
item = self.inbox_rx.recv() => {
match item {
Ok(item) => self.on_inbox(item).await,
Err(broadcast::error::RecvError::Lagged(_)) => {}
Err(broadcast::error::RecvError::Closed) => break,
}
}
attached = self.attach_rx.recv() => {
if let Some(pid) = attached {
self.on_attach(pid);
}
}
state = self.state_rx.recv() => {
match state {
Ok(MindState::Idle) => self.interrupt_at = None,
Ok(_) => {}
Err(broadcast::error::RecvError::Lagged(_)) => {
if self.runtime.mind_state() == MindState::Idle {
self.interrupt_at = None;
}
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
turn = self.turn_rx.recv() => {
if let Ok(turn) = turn {
self.on_turn_frame(&turn.turn);
}
}
_ = sleep_until_opt(self.interrupt_at), if self.interrupt_at.is_some() => {
self.interrupt_at = None;
tracing::error!(
wave = self.runtime.name(),
"interrupt deadline expired with the resident silent; force-finalizing"
);
self.runtime.force_finalize_open_turn(
Lifecycle::Interrupted,
"interrupt deadline: resident silent; listener force-finalized",
);
}
_ = sleep_until_opt(self.respawn_at), if self.respawn_at.is_some() => {
self.respawn_at = None;
self.attempts += 1;
tracing::info!(
wave = self.runtime.name(),
attempt = self.attempts,
"auto-respawning the resident"
);
self.spawn().await;
}
_ = probe.tick() => {
self.probe_attached().await;
}
}
}
}
async fn spawn(&mut self) {
if self.spawner.is_none() {
return;
}
self.respawn_at = None;
if self.child.is_none() {
if let Some(pid) = self.door.seat_pid() {
if process_alive(pid).await {
tracing::info!(
wave = self.runtime.name(),
pid,
"a live resident already holds the seat; not spawning"
);
return;
}
}
}
let spawner = self.spawner.as_mut().expect("checked above");
match spawner() {
Ok(child) => {
if let Some(pid) = child.id() {
self.door.record_pid(pid);
}
self.runtime.set_resident_expected();
self.child = Some(child);
}
Err(err) => {
let reason = format!("resident spawn failed: {err}");
tracing::error!(wave = self.runtime.name(), error = %err, "resident spawn failed");
self.mark_failed(&reason);
if let Some(delay) = respawn_delay(&self.config.respawn_backoff, self.attempts) {
self.respawn_at = Some(Instant::now() + delay);
}
}
}
}
fn on_resident_death(&mut self, reason: &str) {
tracing::error!(wave = self.runtime.name(), reason, "resident died");
self.door.clear_seat();
self.interrupt_at = None;
self.runtime
.force_finalize_open_turn(Lifecycle::Failed, "resident died mid-turn");
self.mark_failed(reason);
if self.spawner.is_some() {
if let Some(delay) = respawn_delay(&self.config.respawn_backoff, self.attempts) {
self.respawn_at = Some(Instant::now() + delay);
}
}
}
fn mark_failed(&self, reason: &str) {
if !matches!(self.runtime.mind_state(), MindState::Failed { .. }) {
self.runtime.transition(
MindState::Failed {
reason: reason.to_string(),
},
reason,
);
}
}
async fn on_inbox(&mut self, item: InboxItem) {
let is_interrupt = match &item {
InboxItem::Interrupt => true,
InboxItem::Message(message) => {
if self.spawner.is_some()
&& self.child.is_none()
&& matches!(self.runtime.mind_state(), MindState::Failed { .. })
{
tracing::info!(
wave = self.runtime.name(),
"message for a dead resident; respawning now"
);
self.spawn().await;
}
message.op == MessageOp::Interrupt
}
};
if is_interrupt
&& self.interrupt_at.is_none()
&& matches!(
self.runtime.mind_state(),
MindState::Turning { .. } | MindState::Interrupting { .. }
)
{
self.interrupt_at = Some(Instant::now() + self.config.interrupt_deadline);
}
}
fn on_attach(&mut self, pid: u32) {
if self
.child
.as_ref()
.is_some_and(|child| child.id() == Some(pid))
{
return;
}
if let Some(mut child) = self.child.take() {
if let Some(old_pid) = child.id() {
tracing::warn!(
wave = self.runtime.name(),
old_pid,
new_pid = pid,
"attach replaced a spawned resident; terminating the old one"
);
tokio::spawn(async move {
terminate_resident(old_pid).await;
let _ = child.wait().await;
});
}
}
self.respawn_at = None;
self.attempts = 0;
tracing::info!(
wave = self.runtime.name(),
pid,
"resident attached; respawn stood down"
);
}
fn on_turn_frame(&mut self, turn: &ChatTurn) {
if turn.role == ChatRole::Assistant && turn.status == Lifecycle::Completed {
self.attempts = 0;
}
}
async fn probe_attached(&mut self) {
if self.child.is_some() {
return;
}
let Some(pid) = self.door.seat_pid() else {
return;
};
if !process_alive(pid).await {
self.on_resident_death("resident pid gone");
}
}
}
pub async fn terminate_resident(pid: u32) {
signal_pid(pid, "-TERM");
for _ in 0..30 {
if !process_alive(pid).await {
return;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
tracing::warn!(pid, "resident ignored SIGTERM; killing");
signal_pid(pid, "-KILL");
}
pub fn terminate_resident_blocking(pid: u32) {
let _ = std::process::Command::new("kill")
.args(["-TERM", &pid.to_string()])
.status();
}
fn signal_pid(pid: u32, signal: &str) {
let _ = std::process::Command::new("kill")
.args([signal, &pid.to_string()])
.status();
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use crate::wave::journal::{journal_path, EventKind, Journal};
use crate::wave::wire::ResidentDelta;
fn open_runtime(repo: &std::path::Path) -> Arc<WaveRuntime> {
WaveRuntime::open("ship".into(), repo.to_path_buf()).expect("open runtime")
}
fn counting_spawner(script: &'static str) -> (SpawnResident, Arc<AtomicU32>) {
let count = Arc::new(AtomicU32::new(0));
let counter = count.clone();
let spawner: SpawnResident = Box::new(move || {
counter.fetch_add(1, Ordering::SeqCst);
tokio::process::Command::new("sh")
.args(["-c", script])
.kill_on_drop(true)
.spawn()
});
(spawner, count)
}
fn config(backoff: Vec<Duration>) -> SupervisorConfig {
SupervisorConfig {
respawn_backoff: backoff,
interrupt_deadline: Duration::from_millis(80),
attach_probe: Duration::from_millis(40),
}
}
async fn wait_for(what: &str, cond: impl Fn() -> bool) {
for _ in 0..500 {
if cond() {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("condition not met in time: {what}");
}
#[test]
fn respawn_delay_walks_the_ladder_and_caps() {
let ladder = RESPAWN_BACKOFF.to_vec();
assert_eq!(respawn_delay(&ladder, 0), Some(Duration::from_secs(300)));
assert_eq!(respawn_delay(&ladder, 1), Some(Duration::from_secs(900)));
assert_eq!(respawn_delay(&ladder, 2), Some(Duration::from_secs(2700)));
assert_eq!(
respawn_delay(&ladder, 9),
Some(Duration::from_secs(2700)),
"past the ladder the cap repeats"
);
assert_eq!(respawn_delay(&[], 0), None, "empty ladder disables");
}
#[tokio::test]
async fn resident_death_fails_the_mind_and_the_ladder_respawns() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let (spawner, spawns) = counting_spawner("exit 0");
let task = tokio::spawn(
Supervisor::new(
rt.clone(),
door,
Some(spawner),
config(vec![Duration::from_millis(30)]),
)
.run(),
);
wait_for("mind failed", || {
matches!(rt.mind_state(), MindState::Failed { .. })
})
.await;
assert!(rt.resident_expected(), "spawn marks the resident expected");
wait_for("respawns happened", || spawns.load(Ordering::SeqCst) >= 3).await;
task.abort();
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("journal");
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::MindState { to: MindState::Failed { reason }, .. }
if reason.contains("resident process exited")
)),
"the death is journaled with its reason"
);
}
#[tokio::test]
async fn resident_death_closes_the_open_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let (spawner, _spawns) = counting_spawner("sleep 0.3");
let task = tokio::spawn(
Supervisor::new(rt.clone(), door, Some(spawner), config(Vec::new())).run(),
);
rt.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
rt.apply_resident_delta(ResidentDelta::TurnText {
text: "half".into(),
});
wait_for("mind failed", || {
matches!(rt.mind_state(), MindState::Failed { .. })
})
.await;
task.abort();
let thread = rt.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].status, Lifecycle::Failed, "open turn closed");
assert_eq!(thread[0].text, "half");
rt.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
});
assert_eq!(rt.thread_snapshot().len(), 1);
assert!(matches!(rt.mind_state(), MindState::Failed { .. }));
}
#[tokio::test]
async fn human_message_respawns_a_dead_resident_immediately() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let (spawner, spawns) = counting_spawner("exit 0");
let task = tokio::spawn(
Supervisor::new(
rt.clone(),
door,
Some(spawner),
config(vec![Duration::from_secs(3600)]),
)
.run(),
);
wait_for("mind failed", || {
matches!(rt.mind_state(), MindState::Failed { .. })
})
.await;
let before = spawns.load(Ordering::SeqCst);
rt.deliver_user_message("are you alive?".into(), MessageOp::Message);
wait_for("immediate respawn", || {
spawns.load(Ordering::SeqCst) > before
})
.await;
task.abort();
}
#[tokio::test]
async fn interrupt_deadline_force_finalizes_when_the_resident_is_silent() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let task = tokio::spawn(Supervisor::new(rt.clone(), door, None, config(Vec::new())).run());
rt.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
rt.apply_resident_delta(ResidentDelta::TurnText {
text: "half".into(),
});
rt.deliver_interrupt();
wait_for("force-finalized interrupted", || {
rt.thread_snapshot()
.last()
.is_some_and(|t| t.status == Lifecycle::Interrupted)
})
.await;
assert_eq!(rt.mind_state(), MindState::Idle);
task.abort();
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("journal");
let fold = crate::wave::journal::fold_thread(&events);
assert!(fold.open.is_empty());
assert_eq!(fold.turns.last().unwrap().status, Lifecycle::Interrupted);
}
#[tokio::test]
async fn attach_disarms_the_respawn_ladder() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let (spawner, spawns) = counting_spawner("exit 0");
let sup = Supervisor::new(
rt.clone(),
door.clone(),
Some(spawner),
config(vec![Duration::from_millis(120)]),
);
let handle = sup.handle();
let task = tokio::spawn(sup.run());
wait_for("mind failed", || {
matches!(rt.mind_state(), MindState::Failed { .. })
})
.await;
let before = spawns.load(Ordering::SeqCst);
door.record_pid(std::process::id());
handle.on_attach(std::process::id());
tokio::time::sleep(Duration::from_millis(400)).await;
assert_eq!(
spawns.load(Ordering::SeqCst),
before,
"no respawn over the attached resident"
);
assert_eq!(
door.seat_pid(),
Some(std::process::id()),
"the attached resident keeps its seat"
);
task.abort();
}
#[tokio::test]
async fn respawn_deadline_never_spawns_over_a_live_seat() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let (spawner, spawns) = counting_spawner("exit 0");
let task = tokio::spawn(
Supervisor::new(
rt.clone(),
door.clone(),
Some(spawner),
config(vec![Duration::from_millis(60)]),
)
.run(),
);
wait_for("mind failed", || {
matches!(rt.mind_state(), MindState::Failed { .. })
})
.await;
let before = spawns.load(Ordering::SeqCst);
door.record_pid(std::process::id());
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(
spawns.load(Ordering::SeqCst),
before,
"the deadline aborted rather than spawning over a live seat"
);
task.abort();
}
#[tokio::test]
async fn attached_resident_death_is_detected_by_pid_probe() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
door.record_pid(4_000_000);
rt.set_resident_expected();
let task =
tokio::spawn(Supervisor::new(rt.clone(), door.clone(), None, config(Vec::new())).run());
wait_for("mind failed via probe", || {
matches!(rt.mind_state(), MindState::Failed { .. })
})
.await;
assert!(door.seat_pid().is_none(), "the seat is freed");
task.abort();
}
#[tokio::test]
async fn completed_turn_resets_the_ladder() {
let tmp = tempfile::tempdir().expect("tempdir");
let rt = open_runtime(tmp.path());
let door = ResidentDoor::new("tok");
let count = Arc::new(AtomicU32::new(0));
let counter = count.clone();
let spawner: SpawnResident = Box::new(move || {
let n = counter.fetch_add(1, Ordering::SeqCst);
let script = if n == 1 { "sleep 0.5" } else { "exit 0" };
tokio::process::Command::new("sh")
.args(["-c", script])
.kill_on_drop(true)
.spawn()
});
let task = tokio::spawn(
Supervisor::new(
rt.clone(),
door,
Some(spawner),
config(vec![Duration::from_millis(30), Duration::from_secs(600)]),
)
.run(),
);
wait_for("second spawn", || count.load(Ordering::SeqCst) >= 2).await;
rt.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
rt.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
});
wait_for("third spawn arrives fast (ladder reset)", || {
count.load(Ordering::SeqCst) >= 3
})
.await;
task.abort();
}
}