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 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);
#[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 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,
}
pub struct Live {
pub id: String,
inner: Mutex<Inner>,
tx: broadcast::Sender<Arc<str>>,
wake: Notify,
}
pub struct Start<'a> {
pub cwd: &'a str,
pub cmd: &'a str,
pub desk: &'a str,
pub slot: i64,
pub cols: u16,
pub rows: u16,
pub accent: &'a str,
}
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)>>,
}
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()),
});
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::default(),
old,
unsaved: false,
run: 0,
cwd: String::new(),
}),
tx,
wake: 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>) {
let live: Vec<Arc<Live>> = self.live.lock().unwrap().values().cloned().collect();
let mut by_dir: HashMap<String, Vec<Arc<Live>>> = HashMap::new();
for l in live {
if let Some(cwd) = l.running_in() {
by_dir.entry(cwd).or_default().push(l);
}
}
self.git
.lock()
.unwrap()
.retain(|d, _| by_dir.contains_key(d));
for (dir, 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);
(crate::project::head_of(p), crate::project::modified(p))
})
.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_default()
}
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);
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 close(&self, id: &str) {
let l = self.live.lock().unwrap().remove(id);
if let Some(l) = l {
l.stop();
}
self.told.lock().unwrap().remove(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.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 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 i = self.inner.lock().unwrap();
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();
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,
};
i.unsaved = true;
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,
..
} = &mut *i;
screen.feed(parser, &buf[..n]);
*unsaved = true;
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 {
tokio::time::sleep(gap - since).await;
}
let Some(l) = me.upgrade() else { return };
let mut i = l.inner.lock().unwrap();
let Inner {
screen,
shown,
status,
old,
unsaved,
..
} = &mut *i;
if screen.scrollback_cleared() && !old.is_empty() {
old.clear();
*unsaved = true;
}
let frame = screen.frame(&l.id, shown);
gap = match &frame {
Some(f) if f.len() > HEAVY_FRAME => SLOW_FRAME,
_ => FRAME,
};
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 running_in(&self) -> Option<String> {
let i = self.inner.lock().unwrap();
i.status
.running
.then(|| i.cwd.clone())
.filter(|c| !c.is_empty())
}
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());
}
}
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 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,
cmd: "printf 'pane=%s\\n' \"$SNYVI_SESSION\"; pwd; exit 3",
desk: "d",
slot: 1,
cols: 400,
rows: 10,
accent: "",
},
&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_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,
cmd: "read x",
desk: "d",
slot: 1,
cols: 80,
rows: 10,
accent: "",
},
&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,
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: "",
},
&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();
}
}