use std::{
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
mpsc::{Receiver, RecvTimeoutError, Sender},
},
time::{Duration, Instant},
};
use crate::{
protocol::{Command, Event},
supervisor::Supervisor,
};
pub enum Wake {
Cmd(Command),
Output,
Hangup,
}
pub type Waker = Arc<Mutex<Option<Sender<Wake>>>>;
pub enum LoopExit {
Shutdown,
ClientGone,
}
const FRAME_MIN: Duration = Duration::from_millis(8);
const FALLBACK: Duration = Duration::from_millis(200);
fn wait_for(dirty: bool, since_last_tick: Duration) -> Duration {
let target = if dirty { FRAME_MIN } else { FALLBACK };
target.saturating_sub(since_last_tick)
}
fn ready_to_tick(dirty: bool, since_last_tick: Duration) -> bool {
since_last_tick >= FRAME_MIN && (dirty || since_last_tick >= FALLBACK)
}
pub fn run_loop(
sup: &mut Supervisor,
wake_rx: &Receiver<Wake>,
stop: &AtomicBool,
mut emit: impl FnMut(&Event) -> bool,
) -> LoopExit {
let mut dirty = false;
let mut last_tick = Instant::now()
.checked_sub(FALLBACK)
.unwrap_or_else(Instant::now);
loop {
if stop.load(Ordering::Relaxed) {
sup.apply(Command::Shutdown);
return LoopExit::Shutdown;
}
match wake_rx.recv_timeout(wait_for(dirty, last_tick.elapsed())) {
Ok(w) => {
if let Some(exit) = apply(sup, w, &mut dirty) {
return exit;
}
}
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => return LoopExit::ClientGone,
}
while let Ok(w) = wake_rx.try_recv() {
if let Some(exit) = apply(sup, w, &mut dirty) {
return exit;
}
}
if ready_to_tick(dirty, last_tick.elapsed()) {
last_tick = Instant::now();
sup.tick();
for ev in sup.drain() {
if !emit(&ev) {
return LoopExit::ClientGone;
}
}
dirty = false;
}
}
}
fn apply(sup: &mut Supervisor, wake: Wake, dirty: &mut bool) -> Option<LoopExit> {
match wake {
Wake::Cmd(Command::Shutdown) => {
sup.apply(Command::Shutdown);
Some(LoopExit::Shutdown)
}
Wake::Cmd(cmd) => {
sup.apply(cmd);
*dirty = true;
None
}
Wake::Output => {
*dirty = true;
None
}
Wake::Hangup => Some(LoopExit::ClientGone),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn frame_floor_gates_ticks_when_dirty() {
assert!(!ready_to_tick(true, Duration::from_millis(3)));
assert_eq!(
wait_for(true, Duration::from_millis(3)),
Duration::from_millis(5)
);
assert!(ready_to_tick(true, Duration::from_millis(8)));
assert_eq!(wait_for(true, Duration::from_millis(8)), Duration::ZERO);
}
#[test]
fn idle_waits_the_backstop_not_the_floor() {
assert!(!ready_to_tick(false, Duration::from_millis(8)));
assert!(!ready_to_tick(false, Duration::from_millis(199)));
assert!(ready_to_tick(false, Duration::from_millis(200)));
assert_eq!(
wait_for(false, Duration::from_millis(50)),
Duration::from_millis(150)
);
}
#[test]
fn watched_input_echoes_without_polling_delay() {
use std::{sync::mpsc::channel, thread};
let cwd = std::env::current_dir().unwrap();
let mut sup = Supervisor::new(24, 80, 2000);
sup.set_launch_context(crate::protocol::LaunchContext::here());
let (wake_tx, wake_rx) = channel::<Wake>();
let (evt_tx, evt_rx) = channel::<Event>();
sup.set_waker(wake_tx.clone());
let core = thread::spawn(move || {
let stop = AtomicBool::new(false);
run_loop(&mut sup, &wake_rx, &stop, |ev| {
evt_tx.send(ev.clone()).is_ok()
});
});
wake_tx
.send(Wake::Cmd(Command::Spawn {
command: "cat".into(),
cwd,
group: None,
}))
.unwrap();
wake_tx
.send(Wake::Cmd(Command::Watch {
id: Some(1),
attached: true,
}))
.unwrap();
let up = wait_for_screen(&evt_rx, |_| true, Duration::from_secs(5));
assert!(up, "watched task never produced an initial screen");
let sent = Instant::now();
wake_tx
.send(Wake::Cmd(Command::Input {
id: 1,
bytes: b"zqmarkerqz\n".to_vec(),
}))
.unwrap();
let echoed = wait_for_screen(
&evt_rx,
|sv| sv.lines.iter().any(|l| l.contains("zqmarkerqz")),
Duration::from_secs(5),
);
let latency = sent.elapsed();
assert!(echoed, "the echo never reached the client");
eprintln!("echo latency: {latency:?}");
assert!(
latency < Duration::from_millis(50),
"echo took {latency:?}: expected an event-driven wake, not a poll"
);
wake_tx.send(Wake::Cmd(Command::Shutdown)).unwrap();
core.join().unwrap();
}
#[test]
fn stop_flag_ends_loop_with_shutdown() {
let cwd = std::env::current_dir().unwrap();
let mut sup = Supervisor::new(24, 80, 2000);
sup.set_launch_context(crate::protocol::LaunchContext::here());
let (wake_tx, wake_rx) = std::sync::mpsc::channel::<Wake>();
sup.set_waker(wake_tx);
sup.apply(Command::Spawn {
command: "sleep 30".into(),
cwd,
group: None,
});
let stop = AtomicBool::new(true);
let started = Instant::now();
let exit = run_loop(&mut sup, &wake_rx, &stop, |_| true);
assert!(matches!(exit, LoopExit::Shutdown));
assert!(started.elapsed() < FALLBACK);
sup.tick();
assert!(
sup.drain()
.iter()
.any(|e| matches!(e, Event::Tasks(v) if v.is_empty()))
);
}
fn wait_for_screen(
evt_rx: &std::sync::mpsc::Receiver<Event>,
pred: impl Fn(&crate::protocol::ScreenView) -> bool,
budget: Duration,
) -> bool {
let deadline = Instant::now() + budget;
while let Some(remaining) = deadline.checked_duration_since(Instant::now()) {
match evt_rx.recv_timeout(remaining) {
Ok(Event::Screen(sv)) if pred(&sv) => return true,
Ok(_) => {}
Err(_) => break,
}
}
false
}
}