use std::{
collections::BTreeMap,
ffi::OsString,
io,
path::{Path, PathBuf},
sync::{Arc, Mutex, OnceLock, 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, WriteRefused},
};
const IDLE_AFTER: Duration = Duration::from_secs(10);
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 MAX_COMMAND_LEN: usize = 64 * 1024;
const SKIP_REASONS: &str = "missing dir, task limit, or command too long";
pub const FLEETCOM_SCROLLBACK: &str = "FLEETCOM_SCROLLBACK";
pub const DEFAULT_SCROLLBACK: usize = 2000;
const MAX_SCROLLBACK: usize = 100_000;
static SCROLLBACK_FLAG: OnceLock<usize> = OnceLock::new();
pub fn set_scrollback_flag(lines: usize) {
let _ = SCROLLBACK_FLAG.set(lines);
}
pub fn scrollback_flag() -> Option<usize> {
SCROLLBACK_FLAG.get().copied()
}
pub fn resolve_scrollback() -> usize {
effective_scrollback(
scrollback_flag(),
std::env::var(FLEETCOM_SCROLLBACK).ok().as_deref(),
)
}
fn effective_scrollback(flag: Option<usize>, env: Option<&str>) -> usize {
flag.or_else(|| env.and_then(|v| v.parse().ok()))
.map_or(DEFAULT_SCROLLBACK, |lines| lines.min(MAX_SCROLLBACK))
}
const KILL_GRACE: Duration = Duration::from_secs(2);
const RECOVERY_DEBOUNCE: Duration = Duration::from_secs(2);
const RECOVERY_CADENCE: Duration = Duration::from_secs(60);
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())))
}
struct Recovery {
enabled: bool,
dirty: bool,
last_mutation: Option<Instant>,
last_cadence: Instant,
last_written: Option<(PathBuf, String)>,
stem: String,
failing: bool,
debounce: Duration,
cadence: Duration,
}
impl Recovery {
fn new() -> Recovery {
Recovery {
enabled: !cfg!(test),
dirty: false,
last_mutation: None,
last_cadence: Instant::now(),
last_written: None,
stem: session::recovery_stem(std::time::SystemTime::now(), std::process::id()),
failing: false,
debounce: RECOVERY_DEBOUNCE,
cadence: RECOVERY_CADENCE,
}
}
}
pub struct Supervisor {
tasks: Vec<Task>,
graveyard: Vec<Task>,
next_id: u64,
rows: u16,
cols: u16,
scrollback: usize,
watched: Option<u64>,
watch_attached: bool,
last_screen: Option<ScreenView>,
launch: Option<LaunchContext>,
events: Vec<Event>,
waker: Waker,
kill_grace: Duration,
capture: BTreeMap<PathBuf, assets::CaptureAssets>,
recovery: Recovery,
}
impl Supervisor {
pub fn new(rows: u16, cols: u16, scrollback: usize) -> Supervisor {
Supervisor {
tasks: Vec::new(),
graveyard: Vec::new(),
next_id: 1,
rows,
cols,
scrollback,
watched: None,
watch_attached: false,
last_screen: None,
launch: None,
events: Vec::new(),
waker: Arc::new(Mutex::new(None)),
kill_grace: KILL_GRACE,
capture: BTreeMap::new(),
recovery: Recovery::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;
}
#[cfg(test)]
pub fn set_recovery_timing(&mut self, debounce: Duration, cadence: Duration) {
self.recovery.enabled = true;
self.recovery.debounce = debounce;
self.recovery.cadence = cadence;
}
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) {
if let Some(old) = self.watched
&& let Some(t) = self.by_id_mut(old)
{
t.scroll_view(ScrollAction::Live);
}
self.watched = None;
self.watch_attached = false;
self.last_screen = None;
}
pub fn apply(&mut self, cmd: Command) {
if matches!(
&cmd,
Command::Spawn { .. }
| Command::Remove { .. }
| Command::Restart { .. }
| Command::SetGroup { .. }
| Command::SetName { .. }
| Command::LoadSession { .. }
| Command::LoadRecovery { .. }
) {
self.recovery.dirty = true;
self.recovery.last_mutation = Some(Instant::now());
}
match cmd {
Command::Spawn {
command,
cwd,
group,
} => self.spawn(&command, cwd, group),
Command::Kill { id } => self.with_task(id, Task::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 } => self.with_task(id, |t| t.tagged = on),
Command::SetGroup { id, group } => {
self.with_task(id, |t| t.group = normalize_group(group))
}
Command::SetName { id, name } => self.with_task(id, |t| 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, attached } => {
if id != self.watched
&& let Some(old) = self.watched
&& let Some(t) = self.by_id_mut(old)
{
t.scroll_view(ScrollAction::Live);
}
if id != self.watched || attached != self.watch_attached {
if let Some(new) = id
&& let Some(t) = self.by_id_mut(new)
{
let _ = t.drain_clipboard();
}
self.last_screen = None;
}
self.watched = id;
self.watch_attached = attached;
}
Command::Input { id, bytes } => self.deliver(id, "input", |t| t.send_input(&bytes)),
Command::Paste { id, bytes } => self.deliver(id, "paste", |t| t.send_paste(&bytes)),
Command::Mouse { id, kind, col, row } => {
self.deliver(id, "mouse input", |t| t.send_mouse(kind, col, row))
}
Command::Key { id, code, mods } => {
self.deliver(id, "key input", |t| t.send_key(code, mods))
}
Command::Scrollback { id, action } => self.with_task(id, |t| t.scroll_view(action)),
Command::SaveSession { name } => self.save_session(&name),
Command::LoadSession { name } => self.load_session(&name),
Command::LoadRecovery { stem } => self.load_recovery(&stem),
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();
t.finalize_preview();
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();
let now = Instant::now();
let forwarding = if self.watch_attached {
self.watched
} else {
None
};
let mut clipboard = None;
let views = self
.tasks
.iter_mut()
.map(|t| {
t.flush_expired_sync();
let stores = t.drain_clipboard();
if forwarding == Some(t.id) {
clipboard = Some((t.id, stores));
}
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, IDLE_AFTER),
preview: t.resolve_preview(now),
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, stores)) = clipboard {
for (kind, text) in stores.stores {
self.events.push(Event::ClipboardCopy { id, kind, text });
}
if let Some(len) = stores.oversized_len {
self.status(format!(
"clipboard copy dropped: {} exceeds the {} limit",
crate::format::bytes(len),
crate::format::bytes(crate::emulator::CLIPBOARD_STORE_MAX_BYTES)
));
}
}
if let Some(id) = self.watched
&& let Some(t) = self.tasks.iter().find(|t| t.id == id)
{
let (formatted, cursor, hide) = t.formatted();
let (wants_mouse, alt_screen, alt_scroll) = t.input_hints();
let sb = t.scroll_offset();
let mut view = ScreenView {
id,
lines: Vec::new(),
formatted,
cursor,
hide_cursor: hide || sb > 0,
wants_mouse,
alt_screen,
alt_scroll,
scrollback: sb,
};
if self.last_screen.as_ref() != Some(&view) {
self.last_screen = Some(view.clone());
view.lines = t.screen_lines();
self.events.push(Event::Screen(view));
}
}
self.maybe_write_recovery(now);
}
fn maybe_write_recovery(&mut self, now: Instant) {
if !self.recovery.enabled {
return;
}
let debounce_due = self.recovery.dirty
&& self
.recovery
.last_mutation
.is_some_and(|t| now.duration_since(t) >= self.recovery.debounce);
let cadence_due = now.duration_since(self.recovery.last_cadence) >= self.recovery.cadence;
if !(debounce_due || cadence_due) {
return;
}
if cadence_due {
self.recovery.last_cadence = now;
}
if self.tasks.is_empty() {
self.recovery.dirty = false;
return;
}
let Some(root) = self.sessions_root() else {
self.recovery.dirty = false;
return;
};
let cfg = self.refreshed_config();
let hash = fnv1a_hex(session::fingerprint_json(&cfg).as_bytes());
let dest = session::recovery_dir(&root).join(format!("{}.json", self.recovery.stem));
if self
.recovery
.last_written
.as_ref()
.is_some_and(|(r, h)| *r == root && *h == hash)
&& std::fs::metadata(&dest).is_ok()
{
self.recovery.dirty = false;
return;
}
let label = session::recovery_label(std::time::SystemTime::now());
match session::save_recovery_in(
&session::recovery_dir(&root),
&self.recovery.stem,
&label,
&cfg,
) {
Ok(_) => {
self.recovery.last_written = Some((root, hash));
self.recovery.failing = false;
}
Err(e) => {
if !self.recovery.failing {
self.recovery.failing = true;
self.status(format!("recovery snapshot failed: {e}"));
}
}
}
self.recovery.dirty = false;
}
pub fn recovery_maintenance(&mut self) {
self.maybe_write_recovery(Instant::now());
}
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 with_task(&mut self, id: u64, f: impl FnOnce(&mut Task)) {
if let Some(t) = self.by_id_mut(id) {
f(t);
}
}
fn deliver(
&mut self,
id: u64,
what: &str,
f: impl FnOnce(&mut Task) -> Result<(), WriteRefused>,
) {
if let Some(r) = self.by_id_mut(id).and_then(|t| f(t).err()) {
self.notice_refused(id, what, r.len);
}
}
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,
self.scrollback,
&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 command.len() > MAX_COMMAND_LEN {
self.status(format!(
"command too long ({} bytes, limit {}), not spawning",
command.len(),
crate::format::bytes(MAX_COMMAND_LEN)
));
return;
}
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 failed: no task {id}"));
return;
};
scrape_now(&mut self.tasks[i]);
if self.tasks[i].finished.is_none() {
self.status("rerun failed: 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();
if let Some(cap) = &old.capture_file {
let _ = std::fs::remove_file(cap);
}
self.graveyard.push(old);
if self.watched == Some(id) {
self.last_screen = None;
}
}
Err(e) => self.status(format!("spawn failed: {e}")),
}
}
fn refreshed_config(&mut self) -> SessionConfig {
for t in &mut self.tasks {
scrape_now(t);
}
self.session_config()
}
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) {
let cfg = self.refreshed_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, recovery) = self
.sessions_root()
.map(|root| {
(
session::list_in(&root),
session::list_recovery_in(&session::recovery_dir(&root)),
)
})
.unwrap_or_default();
self.events.push(Event::Sessions { names, recovery });
}
fn materialize(&mut self, cfg: &SessionConfig) -> Option<(usize, usize, usize)> {
let launch = self.launch_or_refuse()?;
let (mut spawned, mut skipped, mut failed) = (0usize, 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 || entry.cmd.len() > MAX_COMMAND_LEN {
skipped += 1;
continue;
}
match self.admit(
&entry.cmd,
&resolved,
&launch.env,
entry.group.clone(),
entry.name.clone(),
) {
Ok(()) => spawned += 1,
Err(_) => failed += 1,
}
}
}
Some((spawned, skipped, failed))
}
fn load_config(
&mut self,
subject: &str,
load: impl FnOnce(&Path) -> io::Result<SessionConfig>,
) -> Option<SessionConfig> {
let Some(root) = self.sessions_root() else {
self.status("load failed: no config directory available");
return None;
};
match load(&root) {
Ok(c) => Some(c),
Err(e) if e.kind() == io::ErrorKind::NotFound => {
self.status(format!("{subject} not found"));
None
}
Err(e) => {
self.status(format!("{subject} failed to load: {e}"));
None
}
}
}
fn load_session(&mut self, name: &str) {
let Some(cfg) = self.load_config(&format!("session '{name}'"), |root| {
session::load_in(root, name)
}) else {
return;
};
let Some((spawned, skipped, failed)) = self.materialize(&cfg) else {
return;
};
let mut parts = Vec::new();
if spawned > 0 || (skipped == 0 && failed == 0) {
parts.push(format!("{spawned} task(s)"));
}
if skipped > 0 {
parts.push(format!("{skipped} skipped ({SKIP_REASONS})"));
}
if failed > 0 {
parts.push(format!("{failed} failed to spawn"));
}
self.status(format!("loaded '{name}': {}", parts.join(", ")));
}
fn load_recovery(&mut self, stem: &str) {
let Some(cfg) = self.load_config(&format!("recovery snapshot '{stem}'"), |root| {
session::load_recovery_in(&session::recovery_dir(root), stem)
}) else {
return;
};
let Some((_, skipped, failed)) = self.materialize(&cfg) else {
return;
};
let mut msg = String::from("loaded recovery snapshot; save to name it");
if skipped > 0 {
msg.push_str(&format!(", {skipped} skipped ({SKIP_REASONS})"));
}
if failed > 0 {
msg.push_str(&format!(", {failed} failed to spawn"));
}
self.status(msg);
}
}
#[cfg(test)]
#[path = "supervisor_tests.rs"]
mod tests;