use crate::screen::{self, Screen, Shown};
use anyhow::{bail, Context, Result};
use portable_pty::{ChildKiller, CommandBuilder, MasterPty, PtySize};
use serde::Serialize;
use std::collections::HashMap;
use std::io::{Read, Write};
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tokio::sync::{broadcast, Notify};
const FRAME: Duration = Duration::from_millis(16);
const HEAVY_FRAME: usize = 32 * 1024;
const SLOW_FRAME: Duration = Duration::from_millis(33);
const UNWATCHED_FRAME: Duration = Duration::from_secs(1);
const PERSIST_EVERY: Duration = Duration::from_secs(15);
const GIT_EVERY: Duration = Duration::from_secs(3);
const GIT_BACKOFF: u32 = 10;
const GIT_AT_MOST: Duration = Duration::from_secs(60);
pub const BUSY_OUTPUT: Duration = Duration::from_secs(90);
const RESUME_FOR: Duration = Duration::from_secs(5 * 60);
const MARKS_AT_MOST: Duration = Duration::from_secs(24 * 3600);
const WORKING_SILENT: Duration = Duration::from_secs(10 * 60);
const PRINTING_AT_MOST: Duration = Duration::from_secs(2 * 3600);
#[derive(Clone, Debug, Default, Serialize)]
pub struct Status {
pub running: bool,
pub pid: Option<u32>,
pub since: Option<i64>,
pub exit: Option<i32>,
pub blocked: bool,
pub blocked_since: Option<i64>,
pub cmd: String,
pub title: String,
pub accent: String,
pub branch: String,
pub dirty: bool,
pub agent: &'static str,
pub agent_since: Option<i64>,
pub resume: bool,
#[serde(skip_serializing_if = "std::ops::Not::not")]
pub offer: bool,
#[serde(skip_serializing_if = "String::is_empty")]
pub cwd: String,
#[serde(skip_serializing_if = "String::is_empty")]
pub model: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub ctx_pct: Option<u8>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ctx_size: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ctx_in: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ctx_at: Option<i64>,
}
fn clear_context(s: &mut Status) {
s.model.clear();
s.ctx_pct = None;
s.ctx_size = None;
s.ctx_in = None;
s.ctx_at = None;
}
pub const AGENT_STATES: [&str; 3] = ["working", "needs_you", "done"];
struct Proc {
master: Box<dyn MasterPty + Send>,
writer: Box<dyn Write + Send>,
killer: Box<dyn ChildKiller + Send + Sync>,
}
struct Inner {
screen: Screen,
parser: vte::Parser,
shown: Shown,
proc: Option<Proc>,
status: Status,
old: Vec<String>,
unsaved: bool,
run: u64,
cwd: String,
root: String,
wrote: Instant,
printing_since: Instant,
}
impl Inner {
fn frame_now(&mut self, id: &str) -> Option<String> {
if self.screen.scrollback_cleared() && !self.old.is_empty() {
self.old.clear();
self.unsaved = true;
}
self.screen.frame(id, &mut self.shown)
}
}
pub struct Live {
pub id: String,
inner: Mutex<Inner>,
tx: broadcast::Sender<Arc<str>>,
wake: Notify,
watched: Notify,
}
pub struct Start<'a> {
pub cwd: &'a str,
pub root: &'a str,
pub cmd: &'a str,
pub desk: &'a str,
pub slot: i64,
pub cols: u16,
pub rows: u16,
pub accent: &'a str,
pub offer: bool,
}
type CwdSink = Box<dyn Fn(&str, &str) + Send + Sync>;
pub struct Panes {
live: Mutex<HashMap<String, Arc<Live>>>,
dir: PathBuf,
git: Mutex<HashMap<String, Instant>>,
events: broadcast::Sender<String>,
told: Mutex<HashMap<String, (bool, bool, &'static str)>>,
marks: Mutex<Marks>,
cwd_sink: Mutex<Option<CwdSink>>,
}
impl Panes {
pub fn new(data_dir: &std::path::Path, events: broadcast::Sender<String>) -> Arc<Panes> {
let panes = Arc::new(Panes {
live: Mutex::new(HashMap::new()),
dir: data_dir.join("panes"),
git: Mutex::new(HashMap::new()),
events,
told: Mutex::new(HashMap::new()),
marks: Mutex::new(Marks::default()),
cwd_sink: Mutex::new(None),
});
let weak = Arc::downgrade(&panes);
tokio::spawn(async move {
let mut tick = tokio::time::interval(PERSIST_EVERY);
loop {
tick.tick().await;
let Some(p) = weak.upgrade() else { return };
p.persist_all();
}
});
let weak = Arc::downgrade(&panes);
tokio::spawn(async move {
let mut tick = tokio::time::interval(GIT_EVERY);
loop {
tick.tick().await;
let Some(p) = weak.upgrade() else { return };
p.git_tick().await;
}
});
panes
}
pub fn get(self: &Arc<Self>, id: &str) -> Arc<Live> {
let mut live = self.live.lock().unwrap();
if let Some(l) = live.get(id) {
return l.clone();
}
let old = self.read_text(id);
let (tx, _) = broadcast::channel(256);
let l = Arc::new(Live {
id: id.to_string(),
inner: Mutex::new(Inner {
screen: Screen::new(80, 24),
parser: vte::Parser::new(),
shown: Shown::new(80, 24),
proc: None,
status: Status {
resume: self.marked(id),
offer: self.offered(id),
..Status::default()
},
old,
unsaved: false,
run: 0,
cwd: String::new(),
root: String::new(),
wrote: Instant::now(),
printing_since: Instant::now(),
}),
tx,
wake: Notify::new(),
watched: Notify::new(),
});
live.insert(id.to_string(), l.clone());
drop(live);
tokio::spawn(frames(Arc::downgrade(&l), Arc::downgrade(self)));
l
}
async fn git_tick(self: &Arc<Self>) {
self.follow_folders();
let live: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
let mut by_dir: HashMap<String, (bool, Vec<Arc<Live>>)> = HashMap::new();
for l in live {
if let Some((cwd, home)) = l.running_in() {
let e = by_dir.entry(cwd).or_default();
e.0 |= home;
e.1.push(l);
}
}
self.git
.lock()
.unwrap()
.retain(|d, _| by_dir.contains_key(d));
for (dir, (home, panes)) in by_dir {
let now = Instant::now();
if self
.git
.lock()
.unwrap()
.get(&dir)
.is_some_and(|due| now < *due)
{
continue;
}
let d = dir.clone();
let Ok((branch, dirty)) = tokio::task::spawn_blocking(move || {
let p = std::path::Path::new(&d);
let dirty = if home {
crate::project::modified(p)
} else {
None
};
(crate::project::head_of(p), dirty)
})
.await
else {
continue;
};
let cost = now.elapsed();
let next = now + (cost * GIT_BACKOFF).clamp(GIT_EVERY, GIT_AT_MOST);
self.git.lock().unwrap().insert(dir, next);
for l in panes {
l.set_git(branch.clone().unwrap_or_default(), dirty.unwrap_or(false));
}
}
}
pub fn status(&self, id: &str) -> Status {
self.live
.lock()
.unwrap()
.get(id)
.map(|l| l.inner.lock().unwrap().status.clone())
.unwrap_or_else(|| Status {
resume: self.marked(id),
offer: self.offered(id),
..Status::default()
})
}
pub fn busy(&self) -> Vec<String> {
let now = Instant::now();
let live: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
let mut out: Vec<String> = live
.iter()
.filter(|l| {
let i = l.inner.lock().unwrap();
is_busy(
&i.status,
now.saturating_duration_since(i.wrote),
now.saturating_duration_since(i.printing_since),
)
})
.map(|l| l.id.clone())
.collect();
out.sort();
out
}
pub fn with_agent(&self) -> Vec<String> {
let live: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
live.iter()
.filter(|l| {
let s = &l.inner.lock().unwrap().status;
s.running && !s.agent.is_empty()
})
.map(|l| l.id.clone())
.collect()
}
pub fn mark_resume(&self, ids: Vec<String>) {
{
let mut m = self.marks.lock().unwrap();
m.resume = ids.into_iter().collect();
m.resume_until = None;
m.cap = Instant::now() + MARKS_AT_MOST;
}
self.refresh_marks();
}
pub fn arm_marks(&self) {
{
let mut m = self.marks.lock().unwrap();
if m.resume_until.is_some() {
return;
}
m.resume_until = Some(Instant::now() + resume_for());
}
self.refresh_marks();
}
fn refresh_marks(&self) {
let live: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
for l in live {
let (marked, offered) = (self.marked(&l.id), self.offered(&l.id));
let mut i = l.inner.lock().unwrap();
if !i.status.running {
i.status.resume = marked;
i.status.offer = offered;
}
}
}
pub fn unspent(&self) -> (Vec<String>, Vec<String>) {
let now = Instant::now();
let m = self.marks.lock().unwrap();
let mut resume: Vec<String> = m
.resume
.iter()
.filter(|id| m.resume(id, now))
.cloned()
.collect();
let mut offer: Vec<String> = m
.resume
.union(&m.offer)
.filter(|id| m.offer(id, now))
.cloned()
.collect();
resume.sort();
offer.sort();
(resume, offer)
}
pub fn marked(&self, id: &str) -> bool {
self.marks.lock().unwrap().resume(id, Instant::now())
}
fn unmark(&self, id: &str) {
let mut m = self.marks.lock().unwrap();
m.resume.remove(id);
m.offer.remove(id);
}
pub fn mark_offer(&self, ids: Vec<String>) {
{
let mut m = self.marks.lock().unwrap();
m.offer = ids.into_iter().collect();
m.cap = Instant::now() + MARKS_AT_MOST;
}
self.refresh_marks();
}
pub fn offered(&self, id: &str) -> bool {
self.marks.lock().unwrap().offer(id, Instant::now())
}
pub fn on_cwd(&self, sink: CwdSink) {
*self.cwd_sink.lock().unwrap() = Some(sink);
}
fn follow_folders(&self) {
let live: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
for l in live {
let Some(now) = l.shell_cwd() else { continue };
let was = l.inner.lock().unwrap().cwd.clone();
if was == now || same_folder(&was, &now) {
continue;
}
let s = {
let mut i = l.inner.lock().unwrap();
if i.cwd != was || !i.status.running {
continue;
}
i.cwd = now.clone();
i.status.cwd = now.clone();
i.status.clone()
};
let _ = l.tx.send(status_frame(&l.id, &s).into());
if let Some(sink) = self.cwd_sink.lock().unwrap().as_ref() {
sink(&l.id, &now);
}
}
}
pub fn set_agent(&self, id: &str, state: &str) -> bool {
let Some(l) = self.live.lock().unwrap().get(id).cloned() else {
return false;
};
let state = AGENT_STATES.into_iter().find(|s| *s == state).unwrap_or("");
let mut i = l.inner.lock().unwrap();
if !i.status.running {
return false;
}
if i.status.agent == state {
return true;
}
if state == "needs_you" && !i.status.blocked {
i.status.blocked = true;
i.status.blocked_since = Some(crate::store::now());
} else if i.status.agent == "needs_you" && state != "needs_you" {
i.status.blocked = false;
i.status.blocked_since = None;
}
i.status.agent = state;
i.status.agent_since = (!state.is_empty()).then(crate::store::now);
if state.is_empty() {
clear_context(&mut i.status);
}
let s = i.status.clone();
drop(i);
let _ = l.tx.send(status_frame(&l.id, &s).into());
self.changed(&l.id, &s);
true
}
pub fn set_context(
&self,
id: &str,
model: &str,
pct: Option<u8>,
size: Option<u64>,
input: Option<u64>,
) -> bool {
let Some(l) = self.live.lock().unwrap().get(id).cloned() else {
return false;
};
let mut i = l.inner.lock().unwrap();
if !i.status.running {
return false;
}
let st = &i.status;
if st.model == model && st.ctx_pct == pct && st.ctx_size == size {
i.status.ctx_in = input;
return true;
}
i.status.model = model.to_string();
i.status.ctx_pct = pct;
i.status.ctx_size = size;
i.status.ctx_in = input;
i.status.ctx_at = Some(crate::store::now());
let s = i.status.clone();
drop(i);
let _ = l.tx.send(status_frame(&l.id, &s).into());
self.changed(&l.id, &s);
true
}
pub fn is_running(&self, id: &str) -> bool {
self.live
.lock()
.unwrap()
.get(id)
.is_some_and(|l| l.inner.lock().unwrap().status.running)
}
pub fn running(&self) -> usize {
self.live
.lock()
.unwrap()
.values()
.filter(|l| l.inner.lock().unwrap().status.running)
.count()
}
pub fn forget(&self, id: &str) {
let l = self.live.lock().unwrap().remove(id);
if let Some(l) = l {
let text = {
let mut i = l.inner.lock().unwrap();
i.unsaved.then(|| {
i.unsaved = false;
keep_text(&i.old, i.screen.text())
})
};
if let Some(text) = text {
self.write_text(&l.id, &text);
}
l.stop();
}
self.told.lock().unwrap().remove(id);
}
pub fn discard(&self, id: &str) {
self.forget(id);
let _ = std::fs::remove_file(self.text_path(id));
}
pub fn clear(&self) {
let all: Vec<Arc<Live>> = self.live.lock().unwrap().drain().map(|(_, l)| l).collect();
self.told.lock().unwrap().clear();
for l in all {
l.stop();
}
let _ = std::fs::remove_dir_all(&self.dir);
}
pub fn shutdown(&self) {
self.follow_folders();
self.persist_all();
let all: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
for l in all {
l.stop();
}
}
fn persist_all(&self) {
let all: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
for l in all {
let text = {
let mut i = l.inner.lock().unwrap();
if !i.unsaved {
continue;
}
i.unsaved = false;
keep_text(&i.old, i.screen.text())
};
self.write_text(&l.id, &text);
}
}
fn text_path(&self, id: &str) -> PathBuf {
self.dir.join(format!("{id}.txt"))
}
fn read_text(&self, id: &str) -> Vec<String> {
if !valid_id(id) {
return Vec::new();
}
std::fs::read_to_string(self.text_path(id))
.map(|s| s.lines().map(str::to_string).collect())
.unwrap_or_default()
}
fn write_text(&self, id: &str, lines: &[String]) {
if !valid_id(id) {
return;
}
let _ = std::fs::create_dir_all(&self.dir);
let tmp = self.dir.join(format!("{id}.tmp"));
if std::fs::write(&tmp, lines.join("\n")).is_ok() {
let _ = std::fs::rename(&tmp, self.text_path(id));
}
}
fn changed(&self, id: &str, status: &Status) {
let now = (status.running, status.blocked, status.agent);
if self.told.lock().unwrap().insert(id.to_string(), now) == Some(now) {
return;
}
let _ = self.events.send(format!(
"panes\n{}",
serde_json::json!({ "id": id, "running": status.running, "blocked": status.blocked, "agent": status.agent })
));
}
}
pub fn is_busy(s: &Status, since_output: Duration, printing_for: Duration) -> bool {
if !s.running {
return false;
}
match s.agent {
"needs_you" => true,
"working" => since_output < WORKING_SILENT,
_ => since_output < busy_output() && printing_for < PRINTING_AT_MOST,
}
}
struct Marks {
resume: std::collections::HashSet<String>,
offer: std::collections::HashSet<String>,
resume_until: Option<Instant>,
cap: Instant,
}
impl Default for Marks {
fn default() -> Marks {
Marks {
resume: Default::default(),
offer: Default::default(),
resume_until: None,
cap: Instant::now(),
}
}
}
impl Marks {
fn resume(&self, id: &str, now: Instant) -> bool {
now < self.cap && self.resume_until.is_none_or(|t| now < t) && self.resume.contains(id)
}
fn offer(&self, id: &str, now: Instant) -> bool {
now < self.cap
&& (self.offer.contains(id)
|| (self.resume.contains(id) && self.resume_until.is_some_and(|t| now >= t)))
}
}
fn resume_for() -> Duration {
static FOR: std::sync::OnceLock<Duration> = std::sync::OnceLock::new();
*FOR.get_or_init(|| {
std::env::var("SNYVI_RESUME_S")
.ok()
.and_then(|s| s.parse().ok())
.map(Duration::from_secs)
.unwrap_or(RESUME_FOR)
})
}
fn busy_output() -> Duration {
static QUIET: std::sync::OnceLock<Duration> = std::sync::OnceLock::new();
*QUIET.get_or_init(|| {
std::env::var("SNYVI_QUIET_S")
.ok()
.and_then(|s| s.parse().ok())
.map(Duration::from_secs)
.unwrap_or(BUSY_OUTPUT)
})
}
pub fn valid_id(id: &str) -> bool {
id.len() == 32 && id.bytes().all(|b| b.is_ascii_hexdigit())
}
fn keep_text(old: &[String], now: Vec<String>) -> Vec<String> {
let mut all: Vec<String> = old.iter().cloned().chain(now).collect();
let mut bytes: usize = all.iter().map(|l| l.len() + 1).sum();
let mut drop = 0;
while bytes > screen::SCROLLBACK_BYTES && drop < all.len() {
bytes -= all[drop].len() + 1;
drop += 1;
}
all.drain(..drop);
all
}
impl Live {
pub fn attach(&self) -> (Vec<String>, broadcast::Receiver<Arc<str>>) {
let mut i = self.inner.lock().unwrap();
if let Some(f) = i.frame_now(&self.id) {
let _ = self.tx.send(f.into());
}
self.watched.notify_one();
let mut first = vec![status_frame(&self.id, &i.status)];
if !i.old.is_empty() {
first.push(serde_json::json!({ "t": "old", "p": self.id, "lines": i.old }).to_string());
}
first.push(i.screen.snapshot(&self.id, &i.shown));
(first, self.tx.subscribe())
}
pub fn input(&self, bytes: &[u8], panes: &Panes) {
let mut i = self.inner.lock().unwrap();
let Some(p) = i.proc.as_mut() else { return };
let _ = p.writer.write_all(bytes);
let _ = p.writer.flush();
if i.status.blocked {
i.status.blocked = false;
i.status.blocked_since = None;
let s = i.status.clone();
let _ = self.tx.send(status_frame(&self.id, &s).into());
drop(i);
panes.changed(&self.id, &s);
}
}
pub fn resize(&self, cols: u16, rows: u16) {
let mut i = self.inner.lock().unwrap();
let Inner {
screen,
shown,
proc,
..
} = &mut *i;
let (c, r) = (usize::from(cols), usize::from(rows));
if screen.size() == (c.clamp(2, 1000), r.clamp(1, 500)) {
return;
}
screen.resize(c, r, shown);
if let Some(p) = proc.as_ref() {
let (c, r) = screen.size();
let _ = p.master.resize(PtySize {
rows: r as u16,
cols: c as u16,
pixel_width: 0,
pixel_height: 0,
});
}
drop(i);
self.wake.notify_one();
}
pub fn start(self: &Arc<Self>, s: Start, panes: &Arc<Panes>) -> Result<Status> {
if !std::path::Path::new(s.cwd).is_dir() {
bail!("{} is not a folder any more", s.cwd);
}
let mut i = self.inner.lock().unwrap();
if i.status.running {
bail!("already running");
}
let pty = portable_pty::native_pty_system();
let size = PtySize {
rows: s.rows.clamp(1, 500),
cols: s.cols.clamp(2, 1000),
pixel_width: 0,
pixel_height: 0,
};
let pair = pty.openpty(size).context("opening a terminal")?;
let (mut cmd, born) = command(s.cmd, s.accent);
cmd.cwd(s.cwd);
i.cwd = s.cwd.to_string();
i.root = s.root.to_string();
cmd.env("TERM", "xterm-256color");
cmd.env("COLORTERM", "truecolor");
cmd.env("TERM_PROGRAM", "snyvi");
cmd.env("SNYVI_SESSION", &self.id);
cmd.env("SNYVI_DESK", s.desk);
cmd.env("SNYVI_SLOT", s.slot.to_string());
let mut child = pair
.slave
.spawn_command(cmd)
.with_context(|| format!("starting {}", display_cmd(s.cmd)))?;
drop(pair.slave);
let reader = pair
.master
.try_clone_reader()
.context("reading the terminal")?;
let writer = pair.master.take_writer().context("writing the terminal")?;
let killer = child.clone_killer();
let pid = child.process_id();
let (c, r) = (usize::from(size.cols), usize::from(size.rows));
let mut screen = Screen::new(c, r);
let mut shown = Shown::new(c, r);
std::mem::swap(&mut screen, &mut i.screen);
std::mem::swap(&mut shown, &mut i.shown);
let previous = screen.text();
if !previous.is_empty() {
i.old = keep_text(&i.old, previous);
}
i.parser = vte::Parser::new();
i.run += 1;
let run = i.run;
i.proc = Some(Proc {
master: pair.master,
writer,
killer,
});
i.status = Status {
running: true,
pid,
since: Some(crate::store::now()),
exit: None,
blocked: false,
blocked_since: None,
cmd: s.cmd.to_string(),
title: String::new(),
accent: born,
branch: i.status.branch.clone(),
dirty: i.status.dirty,
agent: "",
agent_since: None,
resume: false,
offer: s.offer,
cwd: s.cwd.to_string(),
..Status::default()
};
i.wrote = Instant::now();
i.printing_since = i.wrote;
i.unsaved = true;
panes.unmark(&self.id);
let status = i.status.clone();
let _ = self.tx.send(status_frame(&self.id, &status).into());
if !i.old.is_empty() {
let old = serde_json::json!({ "t": "old", "p": self.id, "lines": i.old });
let _ = self.tx.send(old.to_string().into());
}
drop(i);
panes.changed(&self.id, &status);
self.wake.notify_one();
let me = Arc::downgrade(self);
let (done_tx, done) = std::sync::mpsc::channel::<()>();
std::thread::Builder::new()
.name(format!("pane {}", &self.id[..8]))
.spawn(move || {
read_loop(me, reader, run);
let _ = done_tx.send(());
})
.context("starting the pane's reader")?;
let me = Arc::downgrade(self);
let panes = Arc::downgrade(panes);
std::thread::Builder::new()
.name(format!("wait {}", &self.id[..8]))
.spawn(move || {
let code = child.wait().map(|s| s.exit_code() as i32).unwrap_or(-1);
let _ = done.recv_timeout(Duration::from_millis(500));
let (Some(me), Some(panes)) = (me.upgrade(), panes.upgrade()) else {
return;
};
me.exited(run, code, &panes);
})
.context("starting the pane's waiter")?;
Ok(status)
}
fn exited(&self, run: u64, code: i32, panes: &Panes) {
let (status, text) = {
let mut i = self.inner.lock().unwrap();
if i.run != run {
return;
}
i.proc = None;
i.status.running = false;
i.status.pid = None;
i.status.exit = Some(code);
i.status.blocked = false;
i.status.blocked_since = None;
i.status.agent = "";
i.status.agent_since = None;
i.unsaved = false;
(i.status.clone(), keep_text(&i.old, i.screen.text()))
};
panes.write_text(&self.id, &text);
let _ = self.tx.send(status_frame(&self.id, &status).into());
panes.changed(&self.id, &status);
self.wake.notify_one();
}
pub fn stop(&self) {
let mut i = self.inner.lock().unwrap();
if let Some(mut p) = i.proc.take() {
let _ = p.killer.kill();
drop(p);
}
}
}
fn read_loop(me: std::sync::Weak<Live>, mut reader: Box<dyn Read + Send>, run: u64) {
let mut buf = vec![0u8; 64 * 1024];
loop {
let n = match reader.read(&mut buf) {
Ok(0) | Err(_) => return,
Ok(n) => n,
};
let Some(l) = me.upgrade() else { return };
{
let mut i = l.inner.lock().unwrap();
if i.run != run {
return;
}
let Inner {
screen,
parser,
proc,
unsaved,
wrote,
printing_since,
..
} = &mut *i;
screen.feed(parser, &buf[..n]);
*unsaved = true;
let now = Instant::now();
if now.saturating_duration_since(*wrote) >= busy_output() {
*printing_since = now;
}
*wrote = now;
if !screen.replies.is_empty() {
let replies = std::mem::take(&mut screen.replies);
if let Some(p) = proc.as_mut() {
let _ = p.writer.write_all(&replies);
let _ = p.writer.flush();
}
}
}
l.wake.notify_one();
}
}
async fn frames(me: std::sync::Weak<Live>, panes: std::sync::Weak<Panes>) {
let mut last = Instant::now() - FRAME;
let mut gap = FRAME;
loop {
let Some(l) = me.upgrade() else { return };
let notified = {
let l2 = l.clone();
drop(l);
async move { l2.wake.notified().await }
};
tokio::select! {
_ = notified => {}
_ = tokio::time::sleep(Duration::from_secs(30)) => continue,
}
let since = last.elapsed();
if since < gap {
let Some(l) = me.upgrade() else { return };
let watched = {
let l2 = l.clone();
drop(l);
async move { l2.watched.notified().await }
};
tokio::select! {
_ = tokio::time::sleep(gap - since) => {}
_ = watched => {}
}
}
let Some(l) = me.upgrade() else { return };
let mut i = l.inner.lock().unwrap();
let frame = i.frame_now(&l.id);
gap = match &frame {
_ if l.tx.receiver_count() == 0 => UNWATCHED_FRAME,
Some(f) if f.len() > HEAVY_FRAME => SLOW_FRAME,
_ => FRAME,
};
let Inner { screen, status, .. } = &mut *i;
if let Some(f) = frame {
let _ = l.tx.send(f.into());
}
last = Instant::now();
let mut said = None;
let mut rang = false;
if std::mem::take(&mut screen.bell) && status.running && !status.blocked {
status.blocked = true;
status.blocked_since = Some(crate::store::now());
rang = true;
}
let retitled = screen.title != status.title;
if retitled {
status.title = screen.title.clone();
}
if rang || retitled {
said = Some(status.clone());
}
if let Some(s) = &said {
let _ = l.tx.send(status_frame(&l.id, s).into());
}
drop(i);
if let (true, Some(s), Some(p)) = (rang, said, panes.upgrade()) {
p.changed(&l.id, &s);
}
}
}
impl Live {
fn shell_cwd(&self) -> Option<String> {
let pid = {
let i = self.inner.lock().unwrap();
if !i.status.running {
return None;
}
i.status.pid?
};
folder_of(pid)
}
fn running_in(&self) -> Option<(String, bool)> {
let (cwd, root) = {
let i = self.inner.lock().unwrap();
if !i.status.running || i.cwd.is_empty() {
return None;
}
(i.cwd.clone(), i.root.clone())
};
let real = |p: &str| std::fs::canonicalize(p).unwrap_or_else(|_| p.into());
let home = !root.is_empty() && real(&cwd).starts_with(real(&root));
Some((cwd, home))
}
fn set_git(&self, branch: String, dirty: bool) {
let mut i = self.inner.lock().unwrap();
if i.status.branch == branch && i.status.dirty == dirty {
return;
}
i.status.branch = branch;
i.status.dirty = dirty;
let s = i.status.clone();
drop(i);
let _ = self.tx.send(status_frame(&self.id, &s).into());
}
}
#[cfg(target_os = "linux")]
pub fn folder_of(pid: u32) -> Option<String> {
let p = std::fs::read_link(format!("/proc/{pid}/cwd")).ok()?;
p.is_dir().then(|| p.to_string_lossy().into_owned())
}
#[cfg(target_os = "macos")]
pub fn folder_of(pid: u32) -> Option<String> {
let mut info: libc::proc_vnodepathinfo = unsafe { std::mem::zeroed() };
let size = std::mem::size_of::<libc::proc_vnodepathinfo>() as libc::c_int;
let n = unsafe {
libc::proc_pidinfo(
pid as libc::c_int,
libc::PROC_PIDVNODEPATHINFO,
0,
&mut info as *mut _ as *mut libc::c_void,
size,
)
};
if n != size {
return None;
}
let raw: &[libc::c_char] = unsafe {
std::slice::from_raw_parts(
info.pvi_cdir.vip_path.as_ptr() as *const libc::c_char,
std::mem::size_of_val(&info.pvi_cdir.vip_path),
)
};
let bytes: Vec<u8> = raw
.iter()
.take_while(|c| **c != 0)
.map(|c| *c as u8)
.collect();
let p = String::from_utf8(bytes).ok()?;
std::path::Path::new(&p).is_dir().then_some(p)
}
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
pub fn folder_of(_pid: u32) -> Option<String> {
None
}
fn same_folder(a: &str, b: &str) -> bool {
match (std::fs::canonicalize(a), std::fs::canonicalize(b)) {
(Ok(x), Ok(y)) => x == y,
_ => false,
}
}
fn status_frame(id: &str, s: &Status) -> String {
serde_json::json!({ "t": "status", "p": id, "s": s }).to_string()
}
fn apply(c: &mut CommandBuilder, d: &crate::prompt::Dress) {
for a in &d.args {
c.arg(a);
}
for (k, v) in &d.env {
c.env(k, v);
}
}
fn command(typed: &str, accent: &str) -> (CommandBuilder, String) {
let typed = typed.trim();
#[cfg(unix)]
{
let shell = std::env::var("SHELL")
.ok()
.filter(|s| !s.is_empty())
.unwrap_or_else(|| "/bin/sh".to_string());
let mut c = CommandBuilder::new(&shell);
let dress = if typed.is_empty() {
crate::prompt::dress(std::path::Path::new(&shell), accent)
} else {
None
};
match dress {
Some(d) => {
apply(&mut c, &d);
(c, crate::prompt::effective(accent))
}
None => {
c.arg("-l");
if !typed.is_empty() {
c.arg("-c");
c.arg(typed);
}
(c, String::new())
}
}
}
#[cfg(windows)]
{
if !typed.is_empty() {
let mut c = CommandBuilder::new("cmd.exe");
c.arg("/C");
c.arg(typed);
return (c, String::new());
}
let prog = std::env::var("ComSpec")
.ok()
.filter(|s| !s.is_empty())
.unwrap_or_else(|| "cmd.exe".to_string());
let mut c = CommandBuilder::new(&prog);
match crate::prompt::dress(std::path::Path::new(&prog), accent) {
Some(d) => {
apply(&mut c, &d);
(c, crate::prompt::effective(accent))
}
None => (c, String::new()),
}
}
}
fn display_cmd(typed: &str) -> &str {
if typed.trim().is_empty() {
"the shell"
} else {
typed.trim()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn kept_text_is_cut_from_the_front_to_the_cap() {
let old: Vec<String> = (0..10).map(|i| format!("old {i}")).collect();
let big = "z".repeat(1024);
let now: Vec<String> = (0..3000).map(|_| big.clone()).collect();
let kept = keep_text(&old, now);
let bytes: usize = kept.iter().map(|l| l.len() + 1).sum();
assert!(bytes <= screen::SCROLLBACK_BYTES);
assert_eq!(kept.last().unwrap(), &big, "the newest line is kept");
assert!(
!kept.iter().any(|l| l.starts_with("old")),
"the oldest go first"
);
let kept = keep_text(&old, vec!["new".into()]);
assert_eq!(kept.len(), 11);
assert_eq!(kept[0], "old 0");
assert_eq!(kept[10], "new");
}
#[test]
fn a_restart_waits_for_agents_and_recent_output_but_not_forever() {
let s = |running: bool, agent: &'static str| Status {
running,
agent,
..Status::default()
};
let long_ago = BUSY_OUTPUT + Duration::from_secs(1);
let just_now = Duration::from_secs(1);
let a_while = Duration::from_secs(60);
let silent = WORKING_SILENT + Duration::from_secs(1);
let for_ever = PRINTING_AT_MOST + Duration::from_secs(1);
assert!(is_busy(&s(true, "working"), long_ago, a_while));
assert!(is_busy(&s(true, ""), just_now, a_while));
assert!(is_busy(&s(true, "done"), just_now, a_while));
assert!(!is_busy(&s(true, ""), long_ago, a_while));
assert!(!is_busy(&s(true, "done"), long_ago, a_while));
assert!(is_busy(&s(true, "needs_you"), long_ago, a_while));
assert!(is_busy(&s(true, "needs_you"), silent, for_ever));
assert!(!is_busy(&s(true, "working"), silent, a_while));
assert!(is_busy(&s(true, ""), just_now, PRINTING_AT_MOST - a_while));
assert!(!is_busy(&s(true, ""), just_now, for_ever));
assert!(!is_busy(&s(true, "done"), just_now, for_ever));
assert!(!is_busy(&s(false, "working"), just_now, a_while));
assert!(!is_busy(&s(false, "needs_you"), just_now, a_while));
assert!(!is_busy(&s(false, ""), just_now, a_while));
}
#[test]
fn a_mark_waits_for_a_look_then_lapses_into_an_offer() {
let t0 = Instant::now();
let m = |until: Option<Instant>| Marks {
resume: ["r".to_string()].into(),
offer: ["o".to_string()].into(),
resume_until: until,
cap: t0 + MARKS_AT_MOST,
};
let later = t0 + Duration::from_secs(10 * 60);
assert!(m(None).resume("r", later));
assert!(!m(None).offer("r", later));
assert!(m(None).offer("o", later));
let armed = m(Some(later + RESUME_FOR));
assert!(armed.resume("r", later + Duration::from_secs(60)));
assert!(!armed.resume("r", later + RESUME_FOR));
assert!(armed.offer("r", later + RESUME_FOR));
assert!(!armed.resume("o", later), "an offer is never a resume");
let day = t0 + MARKS_AT_MOST;
assert!(!m(None).resume("r", day));
assert!(!armed.offer("r", day));
assert!(!armed.offer("o", day));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_resume_mark_rides_the_status_until_the_pane_starts_or_it_lapses() {
let dir = crate::store::tempdir::Dir::new("snyvi-pane-mark");
let (events, _) = broadcast::channel(16);
let panes = Panes::new(&dir.path, events);
let early = "0000000000000000000000000000000a";
let late = "0000000000000000000000000000000b";
let never = "0000000000000000000000000000000c";
let woken = panes.get(early);
assert!(!woken.attach().0[0].contains("\"resume\":true"));
panes.mark_resume(vec![early.into(), late.into()]);
assert!(panes.status(early).resume);
assert!(panes.get(late).attach().0[0].contains("\"resume\":true"));
assert!(panes.status(late).resume);
assert!(!panes.status(never).resume);
assert!(panes.busy().is_empty(), "a stopped pane is never busy");
#[cfg(unix)]
{
let cwd = dir.path.to_string_lossy().to_string();
let s = Start {
cwd: &cwd,
root: &cwd,
cmd: "sleep 30",
desk: "d",
slot: 1,
cols: 80,
rows: 10,
accent: "",
offer: false,
};
let status = panes.get(late).start(s, &panes).unwrap();
assert!(!status.resume);
assert!(!panes.status(late).resume);
assert_eq!(panes.busy(), vec![late.to_string()]);
panes.get(late).stop();
}
assert!(panes.status(early).resume, "the other mark is untouched");
assert!(panes.unspent().0.contains(&early.to_string()));
#[cfg(unix)]
assert!(
!panes.unspent().0.contains(&late.to_string()),
"spent by its start"
);
panes.arm_marks();
panes.marks.lock().unwrap().resume_until = Some(Instant::now());
panes.refresh_marks();
assert!(!panes.status(early).resume);
assert!(panes.status(early).offer);
assert!(!panes.marked(early), "the page's own resume is refused now");
assert!(panes.unspent().0.is_empty());
assert!(
panes.unspent().1.contains(&early.to_string()),
"and written back as an offer"
);
assert!(!panes.get(never).attach().0[0].contains("\"resume\":true"));
}
#[test]
fn only_a_pane_id_is_a_file_name() {
assert!(valid_id("0123456789abcdef0123456789abcdef"));
assert!(!valid_id("../../etc/passwd"));
assert!(!valid_id("0123456789abcdef0123456789abcde/"));
assert!(!valid_id(""));
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_pane_runs_a_process_and_its_output_arrives_as_frames() {
let dir = crate::store::tempdir::Dir::new("snyvi-pane");
let (events, _) = broadcast::channel(16);
let panes = Panes::new(&dir.path, events);
let id = "00112233445566778899aabbccddeeff";
let live = panes.get(id);
let (first, mut rx) = live.attach();
assert!(first[0].contains("\"running\":false"));
let cwd = dir.path.to_string_lossy().to_string();
live.start(
Start {
cwd: &cwd,
root: &cwd,
cmd: "printf 'pane=%s\\n' \"$SNYVI_SESSION\"; pwd; exit 3",
desk: "d",
slot: 1,
cols: 400,
rows: 10,
accent: "",
offer: false,
},
&panes,
)
.unwrap();
let mut seen = String::new();
let mut exit = None;
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while exit.is_none() && tokio::time::Instant::now() < deadline {
let Ok(Ok(msg)) = tokio::time::timeout(Duration::from_secs(10), rx.recv()).await else {
break;
};
let v: serde_json::Value = serde_json::from_str(&msg).unwrap();
if v["t"] == "frame" {
seen.push_str(&msg);
}
if v["t"] == "status" && v["s"]["running"] == false {
exit = v["s"]["exit"].as_i64();
}
}
while let Ok(Ok(msg)) = tokio::time::timeout(Duration::from_millis(300), rx.recv()).await {
seen.push_str(&msg);
}
assert!(seen.contains(&format!("pane={id}")), "{seen}");
assert!(seen.contains(cwd.split('/').next_back().unwrap()), "{seen}");
assert_eq!(exit, Some(3));
let text =
std::fs::read_to_string(dir.path.join("panes").join(format!("{id}.txt"))).unwrap();
assert!(text.contains(&format!("pane={id}")));
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_unwatched_pane_is_caught_up_when_a_page_attaches() {
let dir = crate::store::tempdir::Dir::new("snyvi-pane-idle");
let (events, _) = broadcast::channel(16);
let panes = Panes::new(&dir.path, events);
let live = panes.get("ffeeddccbbaa99887766554433221100");
let cwd = dir.path.to_string_lossy().to_string();
let said = dir.path.join("late-said");
let cmd = format!(
"printf 'early\\n'; sleep 0.3; printf 'late\\n'; : > '{}'; sleep 2",
said.display()
);
live.start(
Start {
cwd: &cwd,
root: &cwd,
cmd: &cmd,
desk: "d",
slot: 1,
cols: 80,
rows: 10,
accent: "",
offer: false,
},
&panes,
)
.unwrap();
for _ in 0..100 {
if said.exists() {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(said.exists(), "the shell never printed its second line");
let (first, _rx) = live.attach();
let snap = first.last().unwrap();
assert!(snap.contains("late"), "{snap}");
}
#[cfg(target_os = "linux")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_folder_through_a_link_is_not_a_move() {
let dir = crate::store::tempdir::Dir::new("snyvi-link");
std::fs::create_dir_all(dir.path.join("real")).unwrap();
std::os::unix::fs::symlink(dir.path.join("real"), dir.path.join("link")).unwrap();
let (events, _ev) = broadcast::channel(64);
let panes = Panes::new(&dir.path, events);
let id = "00112233445566778899aabbccddeeff";
let live = panes.get(id);
let link = dir.path.join("link").to_string_lossy().to_string();
live.start(
Start {
cwd: &link,
root: &link,
cmd: "sleep 30",
desk: "d",
slot: 1,
cols: 80,
rows: 10,
accent: "",
offer: false,
},
&panes,
)
.unwrap();
for _ in 0..50 {
if live.shell_cwd().is_some() {
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert!(
live.shell_cwd().is_some_and(|c| c.ends_with("/real")),
"the kernel names it resolved"
);
panes.follow_folders();
assert_eq!(live.inner.lock().unwrap().cwd, link, "not a move");
assert!(
live.running_in().is_some_and(|(_, home)| home),
"and still inside the desk's folder"
);
live.stop();
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_agent_state_is_set_on_a_running_pane_once_per_change() {
let dir = crate::store::tempdir::Dir::new("snyvi-agent");
let (events, mut ev) = broadcast::channel(64);
let panes = Panes::new(&dir.path, events);
let id = "ffeeddccbbaa99887766554433221100";
assert!(!panes.set_agent(id, "working"), "a pane nobody opened");
let live = panes.get(id);
assert!(!panes.set_agent(id, "working"), "a stopped pane");
let (_, mut rx) = live.attach();
let cwd = dir.path.to_string_lossy().to_string();
live.start(
Start {
cwd: &cwd,
root: &cwd,
cmd: "read x",
desk: "d",
slot: 1,
cols: 80,
rows: 10,
accent: "",
offer: false,
},
&panes,
)
.unwrap();
assert!(panes.set_agent(id, "needs_you"));
assert!(
panes.status(id).blocked,
"needs_you is blocked, for the sidebar"
);
assert!(panes.set_agent(id, "needs_you"));
assert!(panes.set_agent(id, "nonsense"), "an unknown word clears it");
let st = panes.status(id);
assert_eq!(st.agent, "");
assert!(!st.blocked, "leaving needs_you unblocks");
assert!(panes.set_agent(id, "done"));
let mut said = Vec::new();
while let Ok(Ok(m)) = tokio::time::timeout(Duration::from_millis(200), rx.recv()).await {
let v: serde_json::Value = serde_json::from_str(&m).unwrap();
if v["t"] == "status" && v["s"]["running"] == true {
said.push(v["s"]["agent"].as_str().unwrap().to_string());
}
}
assert_eq!(said, ["", "needs_you", "", "done"], "one frame per change");
let mut dots = Vec::new();
while let Ok(m) = ev.try_recv() {
dots.push(m);
}
assert!(
dots.iter().any(|d| d.contains("\"agent\":\"needs_you\"")),
"{dots:?}"
);
live.stop();
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
while panes.status(id).running && tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert_eq!(
panes.status(id).agent,
"",
"the agent goes with its process"
);
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_new_title_reaches_the_desk_and_not_the_sidebar() {
let dir = crate::store::tempdir::Dir::new("snyvi-title");
let (events, mut ev) = broadcast::channel(64);
let panes = Panes::new(&dir.path, events);
let id = "00112233445566778899aabbccddeeff";
let live = panes.get(id);
let (_, mut rx) = live.attach();
let cwd = dir.path.to_string_lossy().to_string();
live.start(
Start {
cwd: &cwd,
root: &cwd,
cmd: "for t in a b c d e; do printf '\\033]0;%s\\007' $t; sleep 0.05; done; printf '\\a'; read x",
desk: "d",
slot: 1,
cols: 80,
rows: 10,
accent: "",
offer: false,
},
&panes,
)
.unwrap();
let mut titles = Vec::new();
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
while !panes.status(id).blocked && tokio::time::Instant::now() < deadline {
while let Ok(m) = rx.try_recv() {
let v: serde_json::Value = serde_json::from_str(&m).unwrap();
if v["t"] == "status" {
titles.push(v["s"]["title"].as_str().unwrap_or("").to_string());
}
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert!(panes.status(id).blocked, "the bell rang");
assert!(
titles.iter().any(|t| t == "c"),
"the header hears titles: {titles:?}"
);
let s = panes.status(id);
panes.changed(id, &s);
let mut dots = Vec::new();
while let Ok(m) = ev.try_recv() {
dots.push(m);
}
assert_eq!(
dots.len(),
2,
"the start and the bell, nothing per title: {dots:?}"
);
assert!(dots[1].contains("\"blocked\":true"), "{dots:?}");
live.stop();
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[test]
fn a_shell_that_moved_is_found_where_it_went() {
let dir = crate::store::tempdir::Dir::new("snyvi-pane-cwd");
std::fs::create_dir(dir.path.join("sub")).unwrap();
let mut child = std::process::Command::new("sh")
.args(["-c", "cd sub && exec sleep 5"])
.current_dir(&dir.path)
.spawn()
.unwrap();
let want = std::fs::canonicalize(dir.path.join("sub")).unwrap();
let mut seen = None;
for _ in 0..50 {
seen = folder_of(child.id()).map(std::path::PathBuf::from);
if seen.as_deref() == Some(want.as_path()) {
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
let _ = child.kill();
let _ = child.wait();
assert_eq!(seen.as_deref(), Some(want.as_path()));
assert_eq!(folder_of(u32::MAX), None, "no such process, no folder");
}
}