use std::path::PathBuf;
use std::sync::mpsc::Sender;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use crate::core::{Wake, Waker};
use crate::path;
use crate::protocol::{Command, Event, ScreenView, TaskView};
use crate::session::{self, SessionConfig};
use crate::task::Task;
const IDLE_AFTER: Duration = Duration::from_millis(600);
type LastScreen = (u64, Vec<u8>, (u16, u16), bool);
const MAX_DIM: u16 = 1000;
const MAX_TASKS: usize = 256;
const KILL_GRACE: Duration = Duration::from_secs(2);
pub struct Supervisor {
tasks: Vec<Task>,
next_id: u64,
rows: u16,
cols: u16,
watched: Option<u64>,
last_screen: Option<LastScreen>,
base_dir: PathBuf,
events: Vec<Event>,
waker: Waker,
kill_grace: Duration,
}
impl Supervisor {
pub fn new(rows: u16, cols: u16, base_dir: PathBuf) -> Supervisor {
Supervisor {
tasks: Vec::new(),
next_id: 1,
rows,
cols,
watched: None,
last_screen: None,
base_dir,
events: Vec::new(),
waker: Arc::new(Mutex::new(None)),
kill_grace: KILL_GRACE,
}
}
#[cfg(test)]
pub fn set_kill_grace(&mut self, grace: Duration) {
self.kill_grace = grace;
}
pub fn set_waker(&self, tx: Sender<Wake>) {
if let Ok(mut slot) = self.waker.lock() {
*slot = Some(tx);
}
}
pub fn clear_waker(&self) {
if let Ok(mut slot) = self.waker.lock() {
*slot = None;
}
}
pub fn clear_watch(&mut self) {
self.watched = None;
self.last_screen = None;
}
pub fn apply(&mut self, cmd: Command) {
match cmd {
Command::Spawn { command, cwd } => self.spawn(&command, cwd),
Command::Kill { id } => {
if let Some(t) = self.by_id_mut(id) {
t.terminate();
}
}
Command::Remove { id } => {
if let Some(i) = self.index_of(id) {
self.tasks.remove(i); }
}
Command::Restart { id } => self.restart(id),
Command::Tag { id, on } => {
if let Some(t) = self.by_id_mut(id) {
t.tagged = on;
}
}
Command::Resize { rows, cols } => {
self.rows = rows.clamp(1, MAX_DIM);
self.cols = cols.clamp(1, MAX_DIM);
for t in &mut self.tasks {
let _ = t.resize(self.rows, self.cols);
}
}
Command::Watch { id } => {
if id != self.watched {
self.last_screen = None;
}
self.watched = id;
}
Command::Input { id, bytes } => {
if let Some(t) = self.by_id_mut(id) {
let _ = t.send_input(&bytes);
}
}
Command::SaveSession { name } => self.save_session(&name),
Command::LoadSession { name } => self.load_session(&name),
Command::Shutdown => self.shutdown_all(),
}
}
pub fn reap(&mut self) {
let now = Instant::now();
for t in &mut self.tasks {
let _ = t.poll_exit();
if t.overdue(now, self.kill_grace) {
t.force_kill();
}
}
}
fn shutdown_all(&mut self) {
for t in &mut self.tasks {
t.terminate();
}
let deadline = Instant::now() + self.kill_grace;
while self.tasks.iter().any(|t| t.finished.is_none()) && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(25));
self.reap();
}
self.tasks.clear(); }
pub fn tick(&mut self) {
self.reap();
let now = Instant::now();
let views = self
.tasks
.iter()
.map(|t| TaskView {
id: t.id,
command: t.command.clone(),
cwd: t.cwd.clone(),
tagged: t.tagged,
lifecycle: t.lifecycle(now, IDLE_AFTER),
preview: t.preview(),
started_ago: now.duration_since(t.started),
})
.collect();
self.events.push(Event::Tasks(views));
if let Some(id) = self.watched
&& let Some(t) = self.tasks.iter().find(|t| t.id == id)
{
let (formatted, cursor, hide_cursor) = t.formatted();
let unchanged = matches!(
&self.last_screen,
Some((lid, lf, lc, lh))
if *lid == id && *lf == formatted && *lc == cursor && *lh == hide_cursor
);
if !unchanged {
self.last_screen = Some((id, formatted.clone(), cursor, hide_cursor));
self.events.push(Event::Screen(ScreenView {
id,
lines: t.screen_lines(),
formatted,
cursor,
hide_cursor,
}));
}
}
}
pub fn drain(&mut self) -> Vec<Event> {
std::mem::take(&mut self.events)
}
fn index_of(&self, id: u64) -> Option<usize> {
self.tasks.iter().position(|t| t.id == id)
}
fn by_id_mut(&mut self, id: u64) -> Option<&mut Task> {
self.tasks.iter_mut().find(|t| t.id == id)
}
fn spawn(&mut self, command: &str, cwd: PathBuf) {
if self.tasks.len() >= MAX_TASKS {
self.events.push(Event::Status(format!(
"task limit reached ({MAX_TASKS}), not spawning"
)));
return;
}
match Task::spawn(
self.next_id,
command,
&cwd,
self.rows,
self.cols,
Arc::clone(&self.waker),
) {
Ok(task) => {
self.next_id += 1;
self.tasks.push(task);
}
Err(e) => self
.events
.push(Event::Status(format!("spawn failed: {e}"))),
}
}
fn restart(&mut self, id: u64) {
let Some(i) = self.index_of(id) else {
self.events
.push(Event::Status(format!("rerun: no task {id}")));
return;
};
if self.tasks[i].finished.is_none() {
self.events
.push(Event::Status("rerun: task is still running".into()));
return;
}
match Task::spawn(
id,
&self.tasks[i].command,
&self.tasks[i].cwd,
self.rows,
self.cols,
Arc::clone(&self.waker),
) {
Ok(mut fresh) => {
fresh.tagged = self.tasks[i].tagged;
self.tasks[i] = fresh;
if self.watched == Some(id) {
self.last_screen = None;
}
}
Err(e) => self
.events
.push(Event::Status(format!("spawn failed: {e}"))),
}
}
fn session_config(&self) -> SessionConfig {
let mut order: Vec<usize> = (0..self.tasks.len()).collect();
order.sort_by_key(|&i| self.tasks[i].id);
let mut cfg = SessionConfig::new();
for &i in &order {
let t = &self.tasks[i];
cfg.entry(path::abbreviate(&t.cwd))
.or_default()
.push(t.command.clone());
}
cfg
}
fn save_session(&mut self, name: &str) {
let cfg = self.session_config();
let count: usize = cfg.values().map(Vec::len).sum();
let status = match session::save(name, &cfg) {
Ok(_) => format!("saved '{name}': {count} command(s)"),
Err(e) => format!("save failed: {e}"),
};
self.events.push(Event::Status(status));
}
fn load_session(&mut self, name: &str) {
let cfg = match session::load(name) {
Ok(c) => c,
Err(_) => {
self.events
.push(Event::Status(format!("session '{name}' not found")));
return;
}
};
let (mut spawned, mut skipped) = (0usize, 0usize);
for (dir, cmds) in &cfg {
let resolved = path::resolve(&self.base_dir, dir);
if !resolved.is_dir() {
skipped += cmds.len();
continue;
}
for cmd in cmds {
if self.tasks.len() >= MAX_TASKS {
skipped += 1;
continue;
}
if let Ok(task) = Task::spawn(
self.next_id,
cmd,
&resolved,
self.rows,
self.cols,
Arc::clone(&self.waker),
) {
self.next_id += 1;
self.tasks.push(task);
spawned += 1;
}
}
}
let status = if skipped > 0 {
format!(
"loaded '{name}': {spawned} task(s), {skipped} skipped (missing dir or task limit)"
)
} else {
format!("loaded '{name}': {spawned} task(s)")
};
self.events.push(Event::Status(status));
}
}
#[cfg(test)]
mod tests {
use std::path::Path;
use super::*;
fn here() -> PathBuf {
std::env::current_dir().unwrap()
}
#[test]
fn session_config_groups_by_dir_in_spawn_order() {
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "a".into(),
cwd: here(),
});
s.apply(Command::Spawn {
command: "b".into(),
cwd: PathBuf::from("/tmp"),
});
s.apply(Command::Spawn {
command: "c".into(),
cwd: here(),
});
let cfg = s.session_config();
assert_eq!(
cfg[&path::abbreviate(&here())],
vec!["a".to_string(), "c".to_string()]
);
assert_eq!(cfg["/tmp"], vec!["b".to_string()]);
}
#[test]
fn tick_emits_snapshot_and_watched_screen() {
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "sleep 30".into(),
cwd: here(),
});
s.tick();
let evs = s.drain();
assert_eq!(evs.len(), 1, "only a Tasks snapshot while unwatched");
let id = match &evs[0] {
Event::Tasks(v) => {
assert_eq!(v.len(), 1);
v[0].id
}
_ => panic!("expected a Tasks snapshot"),
};
s.apply(Command::Watch { id: Some(id) });
s.tick();
let evs = s.drain();
assert!(evs.iter().any(|e| matches!(e, Event::Tasks(_))));
assert!(
evs.iter()
.any(|e| matches!(e, Event::Screen(sv) if sv.id == id)),
"watching a task should stream its Screen"
);
}
#[test]
fn watched_screen_not_resent_when_unchanged() {
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "sleep 30".into(),
cwd: here(),
});
let mut id = 0;
for _ in 0..5 {
s.tick();
for e in s.drain() {
if let Event::Tasks(v) = e
&& let Some(t) = v.first()
{
id = t.id;
}
}
std::thread::sleep(Duration::from_millis(20));
}
assert!(id != 0, "task never appeared");
s.apply(Command::Watch { id: Some(id) });
s.tick();
assert!(
s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"first watched tick sends a full screen"
);
s.tick();
assert!(
!s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"unchanged screen must not be resent"
);
}
fn scratch(tag: &str) -> PathBuf {
let d =
std::env::temp_dir().join(format!("fleetcom_sup_test_{tag}_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
fn spawn_ready(s: &mut Supervisor, command: String, cwd: PathBuf, ready: &Path) -> u64 {
s.apply(Command::Spawn { command, cwd });
for _ in 0..200 {
if ready.exists() {
break;
}
std::thread::sleep(Duration::from_millis(25));
}
assert!(ready.exists(), "task never signalled ready");
s.tick();
match s.drain().first() {
Some(Event::Tasks(v)) => v[0].id,
_ => panic!("expected a Tasks snapshot"),
}
}
fn wait_for_lifecycle(
s: &mut Supervisor,
id: u64,
pred: impl Fn(crate::task::Lifecycle) -> bool,
) {
for _ in 0..200 {
s.tick();
for e in s.drain() {
if let Event::Tasks(v) = e
&& let Some(t) = v.iter().find(|t| t.id == id)
&& pred(t.lifecycle)
{
return;
}
}
std::thread::sleep(Duration::from_millis(25));
}
panic!("task {id} never reached the expected lifecycle");
}
#[test]
fn kill_delivers_term_before_kill() {
use crate::task::Lifecycle;
let dir = scratch("term_first");
let (ready, trapped) = (dir.join("ready"), dir.join("trapped"));
let mut s = Supervisor::new(24, 80, here());
let id = spawn_ready(
&mut s,
format!(
"trap 'echo t > {t}; exit 0' TERM; echo r > {r}; while :; do sleep 0.1; done",
t = trapped.display(),
r = ready.display()
),
dir.clone(),
&ready,
);
s.apply(Command::Kill { id });
wait_for_lifecycle(&mut s, id, |l| l == Lifecycle::Ok);
assert!(trapped.exists(), "the TERM trap never ran");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn term_ignoring_task_escalates_to_kill() {
use crate::task::Lifecycle;
let dir = scratch("escalate");
let ready = dir.join("ready");
let mut s = Supervisor::new(24, 80, here());
s.set_kill_grace(Duration::from_millis(150));
let id = spawn_ready(
&mut s,
format!(
"trap '' TERM; echo r > {r}; while :; do sleep 0.1; done",
r = ready.display()
),
dir.clone(),
&ready,
);
s.apply(Command::Kill { id });
wait_for_lifecycle(&mut s, id, |l| l == Lifecycle::Failed);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn shutdown_returns_early_when_jobs_respect_term() {
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "sleep 300".into(),
cwd: here(),
});
let t0 = Instant::now();
s.apply(Command::Shutdown);
assert!(
t0.elapsed() < Duration::from_secs(1),
"shutdown waited the full grace for a TERM-respecting job"
);
s.tick();
assert!(
s.drain()
.iter()
.any(|e| matches!(e, Event::Tasks(v) if v.is_empty()))
);
}
#[test]
fn shutdown_is_bounded_by_grace() {
let dir = scratch("shutdown_bound");
let ready = dir.join("ready");
let mut s = Supervisor::new(24, 80, here());
s.set_kill_grace(Duration::from_millis(200));
spawn_ready(
&mut s,
format!(
"trap '' TERM; echo r > {r}; while :; do sleep 0.1; done",
r = ready.display()
),
dir.clone(),
&ready,
);
let t0 = Instant::now();
s.apply(Command::Shutdown);
let elapsed = t0.elapsed();
assert!(
elapsed < Duration::from_secs(2),
"shutdown took {elapsed:?}: not bounded by the 200 ms grace"
);
s.tick();
assert!(
s.drain()
.iter()
.any(|e| matches!(e, Event::Tasks(v) if v.is_empty()))
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn clear_watch_stops_screen_stream_and_resets_dedup() {
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "sleep 30".into(),
cwd: here(),
});
s.tick();
let id = match s.drain().first() {
Some(Event::Tasks(v)) => v[0].id,
_ => panic!("expected a Tasks snapshot"),
};
s.apply(Command::Watch { id: Some(id) });
s.tick();
assert!(
s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"watching should stream a Screen"
);
s.clear_watch();
s.tick();
assert!(
!s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"a disconnected client's watch must not keep streaming"
);
s.apply(Command::Watch { id: Some(id) });
s.tick();
assert!(
s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"re-watch after clear_watch must resend the full screen"
);
}
#[test]
fn restart_reruns_finished_task_in_place() {
use crate::task::Lifecycle;
let dir = scratch("restart");
let marker = dir.join("marker");
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: format!("echo run >> {}", marker.display()),
cwd: dir.clone(),
});
s.tick();
let id = match s.drain().first() {
Some(Event::Tasks(v)) => v[0].id,
_ => panic!("expected a Tasks snapshot"),
};
s.apply(Command::Tag { id, on: true });
wait_for_lifecycle(&mut s, id, |l| l == Lifecycle::Ok);
s.apply(Command::Restart { id });
wait_for_lifecycle(&mut s, id, |l| l == Lifecycle::Ok);
let runs = std::fs::read_to_string(&marker).unwrap().lines().count();
assert_eq!(runs, 2, "restart must re-execute the command");
s.tick();
let tagged = s
.drain()
.iter()
.any(|e| matches!(e, Event::Tasks(v) if v.iter().any(|t| t.id == id && t.tagged)));
assert!(tagged, "restart must carry the tag over");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn restart_refuses_running_task_and_unknown_id() {
use crate::task::Lifecycle;
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "sleep 30".into(),
cwd: here(),
});
s.tick();
let id = match s.drain().first() {
Some(Event::Tasks(v)) => v[0].id,
_ => panic!("expected a Tasks snapshot"),
};
s.apply(Command::Restart { id });
assert!(
s.drain()
.iter()
.any(|e| matches!(e, Event::Status(m) if m.contains("still running"))),
"a running task must be refused"
);
s.tick();
let alive = s.drain().iter().any(|e| {
matches!(e, Event::Tasks(v) if v.iter().any(
|t| t.id == id && matches!(t.lifecycle, Lifecycle::Active | Lifecycle::Idle)
))
});
assert!(alive, "the refused task must keep running");
s.apply(Command::Restart { id: 999 });
assert!(
s.drain()
.iter()
.any(|e| matches!(e, Event::Status(m) if m.contains("no task"))),
);
}
#[test]
fn restart_watched_task_resends_screen() {
use crate::task::Lifecycle;
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "true".into(),
cwd: here(),
});
s.tick();
let id = match s.drain().first() {
Some(Event::Tasks(v)) => v[0].id,
_ => panic!("expected a Tasks snapshot"),
};
wait_for_lifecycle(&mut s, id, |l| l == Lifecycle::Ok);
s.apply(Command::Watch { id: Some(id) });
s.tick();
assert!(
s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"first watched tick sends a full screen"
);
s.apply(Command::Restart { id });
s.tick();
assert!(
s.drain().iter().any(|e| matches!(e, Event::Screen(_))),
"restart of the watched task must resend the screen"
);
}
#[test]
fn resize_clamps_hostile_dimensions() {
let mut s = Supervisor::new(24, 80, here());
s.apply(Command::Spawn {
command: "sleep 30".into(),
cwd: here(),
});
s.apply(Command::Resize { rows: 0, cols: 0 });
s.tick(); let _ = s.drain();
s.apply(Command::Resize {
rows: u16::MAX,
cols: u16::MAX,
});
s.tick(); let _ = s.drain();
}
}