use std::{
collections::BTreeMap,
ffi::OsString,
io,
path::{Path, PathBuf},
sync::{Arc, Mutex, mpsc::Sender},
time::{Duration, Instant},
};
use crate::{
core::{Wake, Waker},
harness::{self, assets},
path,
protocol::{Command, Event, LaunchContext, ScreenView, ScrollAction, TaskView, env_get},
session::{self, SessionConfig, SessionEntry},
task::Task,
};
const IDLE_AFTER: Duration = Duration::from_millis(600);
const SORT_IDLE_AFTER: Duration = Duration::from_secs(10);
type LastScreen = (u64, Vec<u8>, (u16, u16), bool, (bool, bool, bool), usize);
const MAX_DIM: u16 = 1000;
const MAX_CELLS: u32 = 500_000;
const _: () = assert!(MAX_CELLS >= MAX_DIM as u32);
const MAX_TASKS: usize = 256;
const KILL_GRACE: Duration = Duration::from_secs(2);
const MAX_LABEL_CHARS: usize = 64;
fn normalize_label(label: Option<String>) -> Option<String> {
let label = label?;
let stripped: String = label.chars().filter(|c| !c.is_control()).collect();
let capped: String = stripped.trim().chars().take(MAX_LABEL_CHARS).collect();
if capped.is_empty() {
return None;
}
Some(capped)
}
fn normalize_group(name: Option<String>) -> Option<String> {
normalize_label(name).filter(|g| g != "Unassigned")
}
fn fnv1a_hex(bytes: &[u8]) -> String {
let mut h: u64 = 0xcbf2_9ce4_8422_2325;
for &b in bytes {
h ^= u64::from(b);
h = h.wrapping_mul(0x0000_0100_0000_01b3);
}
format!("{h:016x}")
}
fn current_resume_id(task: &Task) -> Option<String> {
if let Some(id) = &task.scraped_id {
return Some(id.clone());
}
if let (Some(h), Some(path)) = (task.harness, &task.capture_file)
&& let Ok(payload) = std::fs::read_to_string(path)
&& let Some(id) = h.parse_capture(&payload)
{
return Some(id);
}
task.resume_id.clone()
}
fn scrape_now(t: &mut Task) {
let _ = t.poll_exit();
t.scrape_exit_hint();
}
fn harness_home(env: &[(OsString, OsString)], h: &dyn harness::Harness) -> Option<PathBuf> {
let val = |key: &str| env_get(env, key).map(PathBuf::from);
val(h.home_env_var()).or_else(|| Some(val("HOME")?.join(h.home_dot_dir())))
}
pub struct Supervisor {
tasks: Vec<Task>,
graveyard: Vec<Task>,
next_id: u64,
rows: u16,
cols: u16,
watched: Option<u64>,
last_screen: Option<LastScreen>,
launch: Option<LaunchContext>,
events: Vec<Event>,
waker: Waker,
kill_grace: Duration,
capture: BTreeMap<PathBuf, assets::CaptureAssets>,
}
impl Supervisor {
pub fn new(rows: u16, cols: u16) -> Supervisor {
Supervisor {
tasks: Vec::new(),
graveyard: Vec::new(),
next_id: 1,
rows,
cols,
watched: None,
last_screen: None,
launch: None,
events: Vec::new(),
waker: Arc::new(Mutex::new(None)),
kill_grace: KILL_GRACE,
capture: BTreeMap::new(),
}
}
pub fn set_launch_context(&mut self, ctx: LaunchContext) {
self.launch = Some(ctx);
}
#[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,
group,
} => self.spawn(&command, cwd, group),
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) {
let mut t = self.tasks.remove(i);
t.terminate();
if let Some(cap) = &t.capture_file {
let _ = std::fs::remove_file(cap);
}
self.graveyard.push(t);
}
}
Command::Restart { id } => self.rerun(id),
Command::Tag { id, on } => {
if let Some(t) = self.by_id_mut(id) {
t.tagged = on;
}
}
Command::SetGroup { id, group } => {
if let Some(t) = self.by_id_mut(id) {
t.group = normalize_group(group);
}
}
Command::SetName { id, name } => {
if let Some(t) = self.by_id_mut(id) {
t.name = normalize_label(name);
}
}
Command::Resize { rows, cols } => {
self.rows = rows.clamp(1, MAX_DIM);
self.cols = cols.clamp(1, MAX_DIM);
if u32::from(self.rows) * u32::from(self.cols) > MAX_CELLS {
self.cols = (MAX_CELLS / u32::from(self.rows)) as u16;
}
for t in &mut self.tasks {
let _ = t.resize(self.rows, self.cols);
}
}
Command::Watch { id } => {
if id != self.watched {
if let Some(old) = self.watched
&& let Some(t) = self.by_id_mut(old)
{
t.scroll_view(ScrollAction::Live);
}
self.last_screen = None;
}
self.watched = id;
}
Command::Input { id, bytes } => {
let refused = self.by_id_mut(id).and_then(|t| t.send_input(&bytes).err());
if let Some(r) = refused {
self.notice_refused(id, "input", r.len);
}
}
Command::Paste { id, bytes } => {
let refused = self.by_id_mut(id).and_then(|t| t.send_paste(&bytes).err());
if let Some(r) = refused {
self.notice_refused(id, "paste", r.len);
}
}
Command::Mouse { id, kind, col, row } => {
let refused = self
.by_id_mut(id)
.and_then(|t| t.send_mouse(kind, col, row).err());
if let Some(r) = refused {
self.notice_refused(id, "mouse input", r.len);
}
}
Command::Scrollback { id, action } => {
if let Some(t) = self.by_id_mut(id) {
t.scroll_view(action);
}
}
Command::SaveSession { name } => self.save_session(&name),
Command::LoadSession { name } => self.load_session(&name),
Command::ListSessions => self.list_sessions(),
Command::Shutdown => self.shutdown_all(),
}
}
pub fn reap(&mut self) {
let now = Instant::now();
for t in &mut self.tasks {
let _ = t.poll_exit();
t.scrape_exit_hint();
if t.overdue(now, self.kill_grace) {
t.force_kill();
}
}
for t in &mut self.graveyard {
let _ = t.poll_exit();
if t.overdue(now, self.kill_grace) {
t.force_kill();
}
}
self.graveyard.retain_mut(|t| !t.try_collect());
}
fn shutdown_all(&mut self) {
for t in &mut self.tasks {
t.terminate();
}
let deadline = Instant::now() + self.kill_grace;
while !self.swept() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(25));
self.reap();
}
self.tasks.clear(); self.graveyard.clear();
}
fn swept(&mut self) -> bool {
self.graveyard.is_empty() && self.tasks.iter_mut().all(Task::group_gone)
}
pub fn tick(&mut self) {
self.reap();
for t in &self.tasks {
t.flush_expired_sync();
}
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,
group: t.group.clone(),
name: t.name.clone(),
lifecycle: t.lifecycle(now, IDLE_AFTER),
parked: t.parked(now, SORT_IDLE_AFTER),
preview: t.preview(),
started_ago: now.duration_since(t.started),
quiet_ago: t.finished.is_none().then(|| t.quiet_for(now)),
finished_ago: t.finished.map(|f| now.duration_since(f)),
})
.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 hints = t.input_hints();
let sb = t.scroll_offset();
let unchanged = matches!(
&self.last_screen,
Some((lid, lf, lc, lh, lhints, lsb))
if *lid == id && *lf == formatted && *lc == cursor
&& *lh == hide_cursor && *lhints == hints && *lsb == sb
);
if !unchanged {
self.last_screen = Some((id, formatted.clone(), cursor, hide_cursor, hints, sb));
self.events.push(Event::Screen(ScreenView {
id,
lines: t.screen_lines(),
formatted,
cursor,
hide_cursor: hide_cursor || sb > 0,
wants_mouse: hints.0,
alt_screen: hints.1,
alt_scroll: hints.2,
scrollback: sb,
}));
}
}
}
pub fn drain(&mut self) -> Vec<Event> {
std::mem::take(&mut self.events)
}
fn status(&mut self, msg: impl Into<String>) {
self.events.push(Event::Status(msg.into()));
}
fn notice_refused(&mut self, id: u64, what: &str, len: usize) {
self.status(format!(
"task {id} is not reading input; dropped {} {what}",
crate::format::bytes(len)
));
}
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 launch_or_refuse(&mut self) -> Option<LaunchContext> {
if self.launch.is_none() {
self.status("no launch context; reconnect and retry");
}
self.launch.clone()
}
fn launch_env_path(&self, key: &str) -> Option<PathBuf> {
let ctx = self.launch.as_ref()?;
env_get(&ctx.env, key).map(PathBuf::from)
}
fn ensure_capture_assets(&mut self) -> Option<&assets::CaptureAssets> {
let root = self
.launch_env_path(crate::daemon::FLEETCOM_RUNTIME_DIR)
.or_else(|| {
assets::runtime_root(None).map(|base| {
let key = self
.sessions_root()
.map(PathBuf::into_os_string)
.unwrap_or_default();
base.join(fnv1a_hex(key.as_encoded_bytes()))
})
})?;
if let Ok(key) = std::fs::canonicalize(&root)
&& self.capture.contains_key(&key)
{
return self.capture.get(&key);
}
let installed = assets::CaptureAssets::install(&root, std::process::id()).ok()?;
let key = std::fs::canonicalize(&root).unwrap_or(root);
Some(self.capture.entry(key).or_insert(installed))
}
fn spawn_task(
&mut self,
id: u64,
run: u32,
command: &str,
cwd: &Path,
env: &[(OsString, OsString)],
) -> io::Result<Task> {
let mut exec = std::borrow::Cow::Borrowed(command);
let mut env = std::borrow::Cow::Borrowed(env);
let mut meta = None;
if let Some((h, inv)) = harness::detect(command)
&& let Some(paths) = self.ensure_capture_assets().map(|a| a.paths_for(id, run))
{
let home = harness_home(&env, h);
let plan = h.instrument(&inv, &paths, home.as_deref());
exec = format!("{command}{}", plan.args_suffix).into();
env.to_mut().extend(plan.env);
let resume_id = plan.injected_id.or_else(|| inv.known_id());
meta = Some((h, home, paths.capture_file, resume_id));
}
let mut task = Task::spawn(
id,
command,
&exec,
cwd,
self.rows,
self.cols,
&env,
Arc::clone(&self.waker),
)?;
if let Some((h, home, capture_file, resume_id)) = meta {
task.harness = Some(h);
task.harness_home = home;
task.capture_file = Some(capture_file);
task.resume_id = resume_id;
}
task.run = run;
Ok(task)
}
fn admit(
&mut self,
command: &str,
cwd: &Path,
env: &[(OsString, OsString)],
group: Option<String>,
name: Option<String>,
) -> io::Result<()> {
let mut task = self.spawn_task(self.next_id, 0, command, cwd, env)?;
task.group = normalize_group(group);
task.name = normalize_label(name);
self.next_id += 1;
self.tasks.push(task);
Ok(())
}
fn spawn(&mut self, command: &str, cwd: PathBuf, group: Option<String>) {
if self.tasks.len() >= MAX_TASKS {
self.status(format!("task limit reached ({MAX_TASKS}), not spawning"));
return;
}
let Some(launch) = self.launch_or_refuse() else {
return;
};
if let Err(e) = self.admit(command, &cwd, &launch.env, group, None) {
self.status(format!("spawn failed: {e}"));
}
}
fn rerun(&mut self, id: u64) {
let Some(i) = self.index_of(id) else {
self.status(format!("rerun: no task {id}"));
return;
};
scrape_now(&mut self.tasks[i]);
if self.tasks[i].finished.is_none() {
self.status("rerun: task is still running");
return;
}
let Some(launch) = self.launch_or_refuse() else {
return;
};
let (command, cwd) = {
let old = &self.tasks[i];
let command = match (old.harness, current_resume_id(old)) {
(Some(h), Some(rid)) => h.resume_command(&old.command, &rid),
_ => old.command.clone(),
};
(command, old.cwd.clone())
};
let run = self.tasks[i].run + 1;
match self.spawn_task(id, run, &command, &cwd, &launch.env) {
Ok(mut fresh) => {
fresh.tagged = self.tasks[i].tagged;
fresh.group = self.tasks[i].group.clone();
fresh.name = self.tasks[i].name.clone();
let mut old = std::mem::replace(&mut self.tasks[i], fresh);
old.terminate();
self.graveyard.push(old);
if self.watched == Some(id) {
self.last_screen = None;
}
}
Err(e) => self.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(SessionEntry {
cmd: self.recipe_command(t),
group: t.group.clone(),
name: t.name.clone(),
});
}
cfg
}
fn recipe_command(&self, t: &Task) -> String {
let Some(h) = t.harness else {
return t.command.clone();
};
let id = current_resume_id(t)
.or_else(|| h.correlate_fs(&t.cwd, t.spawned_at, t.harness_home.as_deref()));
match id {
Some(id) => h.resume_command(&t.command, &id),
None => t.command.clone(),
}
}
fn sessions_root(&self) -> Option<PathBuf> {
session::sessions_dir(self.launch_env_path(session::FLEETCOM_CONFIG_DIR))
}
fn save_session(&mut self, name: &str) {
for t in &mut self.tasks {
scrape_now(t);
}
let cfg = self.session_config();
let count: usize = cfg.values().map(Vec::len).sum();
let status = match self
.sessions_root()
.map(|root| session::save_in(&root, name, &cfg))
{
Some(Ok(_)) => format!("saved '{name}': {count} command(s)"),
Some(Err(e)) => format!("save failed: {e}"),
None => "save failed: no config directory available".to_string(),
};
self.status(status);
}
fn list_sessions(&mut self) {
let names = self
.sessions_root()
.map(|root| session::list_in(&root))
.unwrap_or_default();
self.events.push(Event::Sessions(names));
}
fn load_session(&mut self, name: &str) {
let cfg = match self
.sessions_root()
.map(|root| session::load_in(&root, name))
{
Some(Ok(c)) => c,
_ => {
self.status(format!("session '{name}' not found"));
return;
}
};
let Some(launch) = self.launch_or_refuse() else {
return;
};
let (mut spawned, mut skipped) = (0usize, 0usize);
for (dir, entries) in &cfg {
let resolved = path::resolve(&launch.cwd, dir);
if !resolved.is_dir() {
skipped += entries.len();
continue;
}
for entry in entries {
if self.tasks.len() >= MAX_TASKS {
skipped += 1;
continue;
}
if self
.admit(
&entry.cmd,
&resolved,
&launch.env,
entry.group.clone(),
entry.name.clone(),
)
.is_ok()
{
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.status(status);
}
}
#[cfg(test)]
#[path = "supervisor_tests.rs"]
mod tests;