use crate::screen::{self, Line, 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 SETTLE: Duration = Duration::from_millis(4);
const ECHO_WINDOW: Duration = Duration::from_millis(50);
const ECHO_SETTLE: Duration = Duration::from_millis(1);
const SYNC_AT_MOST: Duration = Duration::from_millis(150);
const PERSIST_EVERY: Duration = Duration::from_secs(15);
const PERSIST_SCREEN_EVERY: Duration = Duration::from_secs(60);
const GIT_EVERY: Duration = Duration::from_secs(3);
const GIT_FLOOR: Duration = Duration::from_secs(10);
const GIT_BACKOFF: u32 = 10;
const GIT_AT_MOST: Duration = Duration::from_secs(60);
const GIT_QUIET: Duration = Duration::from_secs(60);
static GIT_RUNS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
pub fn git_runs() -> u64 {
GIT_RUNS.load(std::sync::atomic::Ordering::Relaxed)
}
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, Copy, Debug)]
struct Asked {
at: Instant,
cost: Duration,
watchers: usize,
}
fn status_due(
now: Instant,
last: Option<Asked>,
printed: Option<Instant>,
watchers: usize,
) -> bool {
let Some(a) = last else { return true };
if now < a.at + (a.cost * GIT_BACKOFF).clamp(GIT_FLOOR, GIT_AT_MOST) {
return false;
}
watchers > a.watchers
|| printed.is_some_and(|p| p > a.at)
|| (watchers > 0 && now >= a.at + GIT_QUIET)
}
#[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>,
#[serde(skip_serializing_if = "std::ops::Not::not")]
pub agent_in: bool,
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_used: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ctx_at: Option<i64>,
#[serde(skip)]
pub told_at: i64,
}
fn clear_context(s: &mut Status) {
s.model.clear();
s.ctx_pct = None;
s.ctx_size = None;
s.ctx_used = None;
s.ctx_at = None;
}
fn ctx_figure(used: Option<u64>) -> Option<u64> {
used.map(|n| match n {
0..1_000_000 => n / 1000,
_ => 1_000_000 + n / 100_000,
})
}
pub const AGENT_STATES: [&str; 3] = ["working", "needs_you", "done"];
struct Proc {
master: Box<dyn MasterPty + Send>,
input: std::sync::mpsc::Sender<Vec<u8>>,
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,
unsaved_lines: bool,
saved_at: Instant,
filed: bool,
kept: usize,
txt_bytes: usize,
clears: u64,
run: u64,
cwd: String,
started: Instant,
arrived: bool,
root: String,
wrote: Instant,
printing_since: Instant,
typed: Option<Instant>,
}
impl Inner {
fn fresh_screen(&mut self, c: usize, r: usize) {
let mut screen = Screen::new(c, r);
let mut shown = Shown::new(c, r);
std::mem::swap(&mut screen, &mut self.screen);
std::mem::swap(&mut shown, &mut self.shown);
let previous = screen.text();
if !previous.is_empty() {
self.old = keep_text(&self.old, previous);
}
self.filed = false;
self.clears = 0;
self.parser = vte::Parser::new();
}
fn echo_due(&mut self) -> bool {
match self.typed {
Some(t) if self.wrote >= t => {
self.typed = None;
t.elapsed() < ECHO_WINDOW
}
_ => false,
}
}
}
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.unsaved_lines = true;
self.filed = false;
}
self.screen.frame(id, &mut self.shown)
}
fn take_text(&mut self) -> Text {
let screen = self.screen.screen_text();
let screen_bytes = bytes(screen.iter().map(String::len));
if self.screen.clears() != self.clears {
self.old.clear();
self.filed = false;
}
let since = if self.filed {
self.screen.kept_since(self.kept)
} else {
None
};
if let Some(lines) = since {
let more = bytes(lines.iter().map(Line::text_len));
if self.txt_bytes + more + screen_bytes <= screen::SCROLLBACK_BYTES {
let first = self.txt_bytes == 0;
self.kept = self.screen.lines_ever();
self.txt_bytes += more;
return Text::More {
lines,
first,
screen,
};
}
}
let text = self.screen.text();
let kept = self.screen.kept_lines();
let present = text.len().min(kept);
let shown = text.len() - present;
let all = keep_text(&self.old, text);
let split = all.len().saturating_sub(shown);
let (lines, screen) = all.split_at(split);
self.kept = self.screen.lines_ever() - (kept - present);
self.txt_bytes = bytes(lines.iter().map(String::len));
self.filed = true;
self.clears = self.screen.clears();
Text::Whole {
lines: lines.to_vec(),
screen: screen.to_vec(),
}
}
}
enum Text {
Whole {
lines: Vec<String>,
screen: Vec<String>,
},
More {
lines: Vec<Line>,
first: bool,
screen: Vec<String>,
},
}
fn bytes(lens: impl Iterator<Item = usize>) -> usize {
lens.map(|n| n + 1).sum()
}
pub struct Live {
pub id: String,
inner: Mutex<Inner>,
tx: broadcast::Sender<Arc<str>>,
wake: Notify,
watched: Notify,
slow: std::sync::atomic::AtomicUsize,
}
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,
pub env: &'a [(String, String)],
}
type CwdSink = Box<dyn Fn(&str, &str) + Send + Sync>;
pub struct Panes {
live: Mutex<HashMap<String, Arc<Live>>>,
dir: PathBuf,
git: Mutex<HashMap<String, Asked>>,
events: broadcast::Sender<String>,
writing: Mutex<()>,
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,
writing: Mutex::new(()),
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 };
let _ = tokio::task::spawn_blocking(move || p.persist_all(false)).await;
}
});
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(64);
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,
unsaved_lines: false,
saved_at: Instant::now(),
filed: false,
kept: 0,
txt_bytes: 0,
clears: 0,
run: 0,
cwd: String::new(),
started: Instant::now(),
arrived: false,
root: String::new(),
wrote: Instant::now(),
printing_since: Instant::now(),
typed: None,
}),
tx,
wake: Notify::new(),
watched: Notify::new(),
slow: std::sync::atomic::AtomicUsize::new(0),
});
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();
#[derive(Default)]
struct Folder {
home: bool,
panes: Vec<Arc<Live>>,
watchers: usize,
printed: Option<Instant>,
}
let mut by_dir: HashMap<String, Folder> = HashMap::new();
for l in live {
let watchers = l.tx.receiver_count();
if watchers == 0 {
continue;
}
if let Some((cwd, home)) = l.running_in() {
let wrote = l.inner.lock().unwrap().wrote;
let e = by_dir.entry(cwd).or_default();
e.home |= home;
e.watchers += watchers;
e.printed = e.printed.max(Some(wrote));
e.panes.push(l);
}
}
self.git
.lock()
.unwrap()
.retain(|d, _| by_dir.contains_key(d));
for (dir, f) in by_dir {
let now = Instant::now();
let last = self.git.lock().unwrap().get(&dir).copied();
let ask = f.home && status_due(now, last, f.printed, f.watchers);
let d = dir.clone();
let Ok((branch, dirty)) = tokio::task::spawn_blocking(move || {
let p = std::path::Path::new(&d);
let dirty = if ask {
GIT_RUNS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
crate::project::modified(p)
} else {
None
};
(crate::project::head_of(p), dirty)
})
.await
else {
continue;
};
if ask {
let asked = Asked {
at: now,
cost: now.elapsed(),
watchers: f.watchers,
};
self.git.lock().unwrap().insert(dir, asked);
}
let dirty = if ask || !f.home {
Some(dirty.unwrap_or(false))
} else {
None
};
for l in f.panes {
l.set_git(branch.clone().unwrap_or_default(), dirty);
}
}
}
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) {
const ARRIVE_WITHIN: Duration = Duration::from_secs(1);
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, settling) = {
let i = l.inner.lock().unwrap();
(
i.cwd.clone(),
!i.arrived && i.started.elapsed() < ARRIVE_WITHIN,
)
};
if was == now || same_folder(&was, &now) {
l.inner.lock().unwrap().arrived = true;
continue;
}
if settling {
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 && i.status.agent_in == !state.is_empty() {
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);
i.status.agent_in = !state.is_empty();
if state.is_empty() {
clear_context(&mut i.status);
i.status.told_at = 0;
}
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 agent_in(&self, id: &str) -> 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;
}
if !i.status.agent_in {
i.status.agent_in = true;
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 told(&self, id: &str, now: i64) -> Option<i64> {
let l = self.live.lock().unwrap().get(id).cloned()?;
let mut i = l.inner.lock().unwrap();
if !i.status.running {
return None;
}
let was = i.status.told_at;
i.status.told_at = now;
Some(was)
}
pub fn set_context(
&self,
id: &str,
model: &str,
pct: Option<u8>,
size: Option<u64>,
used: Option<u64>,
) -> Option<bool> {
let l = self.live.lock().unwrap().get(id).cloned()?;
let mut i = l.inner.lock().unwrap();
if !i.status.running {
return None;
}
let st = &i.status;
if st.model == model
&& st.ctx_pct == pct
&& st.ctx_size == size
&& ctx_figure(st.ctx_used) == ctx_figure(used)
{
i.status.ctx_used = used;
return Some(false);
}
i.status.model = model.to_string();
i.status.ctx_pct = pct;
i.status.ctx_size = size;
i.status.ctx_used = used;
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);
Some(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 _w = self.writing.lock().unwrap();
let text = {
let mut i = l.inner.lock().unwrap();
i.unsaved.then(|| {
i.unsaved = false;
i.unsaved_lines = false;
i.take_text()
})
};
if let Some(text) = text {
self.write(&l.id, text);
}
l.stop();
}
self.told.lock().unwrap().remove(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 _w = self.writing.lock().unwrap();
let _ = std::fs::remove_dir_all(&self.dir);
}
pub fn shutdown(&self) {
self.follow_folders();
self.persist_all(true);
let all: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
for l in all {
l.stop();
}
}
fn persist_all(&self, all: bool) {
let panes: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
for l in panes {
let _w = self.writing.lock().unwrap();
if !self.holds(&l) {
continue;
}
let text = {
let mut i = l.inner.lock().unwrap();
let due = all || i.unsaved_lines || i.saved_at.elapsed() >= PERSIST_SCREEN_EVERY;
if !i.unsaved || !due {
continue;
}
i.unsaved = false;
i.unsaved_lines = false;
i.saved_at = Instant::now();
i.take_text()
};
self.write(&l.id, text);
}
}
fn holds(&self, l: &Live) -> bool {
self.live
.lock()
.unwrap()
.get(&l.id)
.is_some_and(|m| std::ptr::eq(Arc::as_ptr(m), l))
}
fn text_path(&self, id: &str) -> PathBuf {
self.dir.join(format!("{id}.txt"))
}
fn screen_path(&self, id: &str) -> PathBuf {
self.dir.join(format!("{id}.scr"))
}
fn read_text(&self, id: &str) -> Vec<String> {
if !valid_id(id) {
return Vec::new();
}
let read = |p: PathBuf| {
std::fs::read_to_string(p)
.map(|s| s.lines().map(str::to_string).collect::<Vec<_>>())
.unwrap_or_default()
};
let mut lines = read(self.text_path(id));
lines.extend(read(self.screen_path(id)));
while lines.last().is_some_and(|l| l.is_empty()) {
lines.pop();
}
keep_text(&[], lines)
}
fn write(&self, id: &str, t: Text) {
if !valid_id(id) {
return;
}
let _ = std::fs::create_dir_all(&self.dir);
let whole = |path: PathBuf, text: String| {
let tmp = path.with_extension("tmp");
if std::fs::write(&tmp, text).is_ok() {
let _ = std::fs::rename(&tmp, path);
}
};
match t {
Text::Whole { lines, screen } => {
whole(self.text_path(id), lines.join("\n"));
whole(self.screen_path(id), screen.join("\n"));
}
Text::More {
lines,
first,
screen,
} => {
if !lines.is_empty() {
let mut more = String::new();
for (n, l) in lines.iter().enumerate() {
if n > 0 || !first {
more.push('\n');
}
more.push_str(&l.text());
}
use std::io::Write as _;
if let Ok(mut f) = std::fs::OpenOptions::new()
.append(true)
.create(true)
.open(self.text_path(id))
{
let _ = f.write_all(more.as_bytes());
}
}
whole(self.screen_path(id), screen.join("\n"));
}
}
}
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())
}
const OLD_BYTES: usize = 64 * 1024;
fn old_frame(id: &str, old: &[String], run: u64, have: Option<usize>, max: usize) -> String {
let end = old.len().saturating_sub(have.unwrap_or(0));
let (mut from, mut size) = (end, 0);
while from > 0 && end - from < max {
let n = old[from - 1].len() + 3;
if size + n > OLD_BYTES && from < end {
break;
}
size += n;
from -= 1;
}
let mut out = String::with_capacity(size + 96);
out.push_str("{\"t\":\"old\",\"p\":");
screen::push_json_str(&mut out, id);
out.push_str(&format!(",\"g\":{run},\"more\":{from}"));
if let Some(h) = have {
out.push_str(&format!(",\"have\":{h}"));
}
out.push_str(",\"lines\":[");
for (k, l) in old[from..end].iter().enumerate() {
if k > 0 {
out.push(',');
}
screen::push_json_str(&mut out, l);
}
out.push_str("]}");
out
}
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(old_frame(&self.id, &i.old, i.run, None, usize::MAX));
}
first.push(i.screen.snapshot(&self.id, &i.shown));
(first, self.tx.subscribe())
}
pub fn pace(&self, slow: bool) {
use std::sync::atomic::Ordering;
if slow {
self.slow.fetch_add(1, Ordering::Relaxed);
} else {
let mut n = self.slow.load(Ordering::Relaxed);
while n > 0 {
match self.slow.compare_exchange_weak(
n,
n - 1,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => break,
Err(now) => n = now,
}
}
self.watched.notify_one();
}
}
pub fn more(&self, old: Option<u64>, before: usize, max: usize) -> Option<String> {
let i = self.inner.lock().unwrap();
match old {
None => Some(i.screen.more(&self.id, before, max)),
Some(run) if run == i.run => {
Some(old_frame(&self.id, &i.old, i.run, Some(before), max))
}
Some(_) => None,
}
}
pub fn input(&self, bytes: &[u8], panes: &Panes) {
let mut i = self.inner.lock().unwrap();
let Some(p) = i.proc.as_ref() else { return };
let _ = p.input.send(bytes.to_vec());
i.typed = Some(Instant::now());
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(dunce::simplified(std::path::Path::new(s.cwd)));
i.cwd = s.cwd.to_string();
i.started = Instant::now();
i.arrived = false;
i.root = s.root.to_string();
fresh_env(&mut cmd, std::env::var_os("CLAUDECODE").is_some());
cmd.env("SNYVI_SESSION", &self.id);
cmd.env("SNYVI_DESK", s.desk);
cmd.env("SNYVI_SLOT", s.slot.to_string());
for (k, v) in s.env {
cmd.env(k, v);
}
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();
i.fresh_screen(usize::from(size.cols), usize::from(size.rows));
i.run += 1;
let run = i.run;
i.proc = Some(Proc {
master: pair.master,
input: write_loop(writer, &self.id),
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;
i.unsaved_lines = 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 = old_frame(&self.id, &i.old, i.run, None, usize::MAX);
let _ = self.tx.send(old.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.status.agent_in = false;
i.unsaved = false;
i.unsaved_lines = false;
i.saved_at = Instant::now();
(i.status.clone(), i.take_text())
};
{
let _w = panes.writing.lock().unwrap();
if panes.holds(self) {
panes.write(&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 write_loop(mut writer: Box<dyn Write + Send>, id: &str) -> std::sync::mpsc::Sender<Vec<u8>> {
let (tx, rx) = std::sync::mpsc::channel::<Vec<u8>>();
let _ = std::thread::Builder::new()
.name(format!("input {}", &id[..8.min(id.len())]))
.spawn(move || {
for bytes in rx {
if writer
.write_all(&bytes)
.and_then(|()| writer.flush())
.is_err()
{
return;
}
}
});
tx
}
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,
unsaved_lines,
wrote,
printing_since,
..
} = &mut *i;
let lines = screen.lines_ever();
screen.feed(parser, &buf[..n]);
*unsaved = true;
if screen.lines_ever() != lines || screen.scrollback_cleared() {
*unsaved_lines = 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_ref() {
let _ = p.input.send(replies);
}
}
}
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;
let mut hold: Option<Instant> = None;
loop {
let Some(l) = me.upgrade() else { return };
let notified = {
let l2 = l.clone();
drop(l);
async move { l2.wake.notified().await }
};
let wait = hold.map_or(Duration::from_secs(30), |t| {
t.saturating_duration_since(Instant::now())
});
tokio::select! {
_ = notified => {}
_ = tokio::time::sleep(wait) => if hold.is_none() { continue },
}
let since = last.elapsed();
let echo = match me.upgrade() {
Some(l) => l.inner.lock().unwrap().echo_due(),
None => return,
};
if echo && hold.is_none() {
tokio::time::sleep(ECHO_SETTLE).await;
} else if since >= gap && hold.is_none() {
tokio::time::sleep(SETTLE).await;
} else 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();
hold = i.screen.holding(SYNC_AT_MOST);
if hold.is_some() {
continue;
}
let frame = i.frame_now(&l.id);
gap = match &frame {
_ if l.tx.receiver_count() <= l.slow.load(std::sync::atomic::Ordering::Relaxed) => {
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: Option<bool>) {
let mut i = self.inner.lock().unwrap();
let dirty = dirty.unwrap_or(i.status.dirty);
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);
}
}
const AGENT_ENV: &[&str] = &[
"CLAUDECODE",
"CLAUDE_CODE_CHILD_SESSION",
"CLAUDE_CODE_SESSION_ID",
"CLAUDE_CODE_SESSION_ATTENDED",
"CLAUDE_CODE_ENTRYPOINT",
"CLAUDE_CODE_EXECPATH",
"CLAUDE_CODE_MESSAGING_SOCKET",
"CLAUDE_CODE_MESSAGING_TOKEN",
"CLAUDE_PID",
"CLAUDE_EFFORT",
"AI_AGENT",
];
const AGENT_TOOL_ENV: &[&str] = &["NO_COLOR", "FORCE_COLOR", "GIT_TERMINAL_PROMPT"];
fn fresh_env(cmd: &mut CommandBuilder, inside: bool) {
for k in AGENT_ENV {
cmd.env_remove(k);
}
if inside {
for k in AGENT_TOOL_ENV {
cmd.env_remove(k);
}
}
cmd.env("TERM", "xterm-256color");
cmd.env("COLORTERM", "truecolor");
cmd.env("TERM_PROGRAM", "snyvi");
}
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;