use std::io::{Read, Write};
use std::sync::{
atomic::{AtomicBool, AtomicU64, Ordering},
mpsc, Arc, Mutex,
};
use std::time::{Duration, Instant};
use anyhow::Result;
use portable_pty::{native_pty_system, CommandBuilder, MasterPty, PtySize};
pub type TermId = usize;
pub enum TermEvent {
Output(TermId),
Exited(TermId),
}
static NEXT_STARTUP_PROBE: AtomicU64 = AtomicU64::new(1);
struct StartupInput {
bytes: Vec<u8>,
marker: String,
probe_interval: Duration,
last_probe: Option<Instant>,
}
pub struct Term {
pub exited: bool,
pub spawn_cwd: Option<std::path::PathBuf>,
parser: Arc<Mutex<vt100::Parser>>,
writer: Box<dyn Write + Send>,
master: Box<dyn MasterPty + Send>,
kill_tx: mpsc::Sender<()>,
startup_input: Option<StartupInput>,
exit_code: Arc<Mutex<Option<i32>>>,
notify_exit: Arc<AtomicBool>,
rows: u16,
cols: u16,
view_offset: usize,
scrollback_limit: usize,
}
pub fn spawn(
id: TermId,
rows: u16,
cols: u16,
scrollback: usize,
cwd: Option<std::path::PathBuf>,
session: Option<&str>,
session_instance_id: Option<&str>,
startup_probe_interval: Duration,
tx: mpsc::Sender<TermEvent>,
) -> Result<Term> {
let rows = rows.max(1);
let cols = cols.max(1);
let pty_system = native_pty_system();
let pair = pty_system.openpty(PtySize {
rows,
cols,
pixel_width: 0,
pixel_height: 0,
})?;
let shell = crate::sys::shell::default_shell();
let mut cmd = CommandBuilder::new(shell);
let spawn_cwd = cwd.filter(|d| d.is_dir());
if let Some(dir) = &spawn_cwd {
cmd.cwd(dir);
}
if let Some(name) = session {
cmd.env("MARS_SESSION", name);
}
if let Some(id) = session_instance_id {
cmd.env("MARS_SESSION_ID", id);
}
let mut child = pair.slave.spawn_command(cmd)?;
drop(pair.slave);
let mut reader = pair.master.try_clone_reader()?;
let writer = pair.master.take_writer()?;
let parser = Arc::new(Mutex::new(vt100::Parser::new(rows, cols, scrollback)));
let reader_parser = parser.clone();
let output_tx = tx.clone();
let ledger_session = session.map(|s| s.to_string());
let ledger_surface = id.to_string();
let (reader_done_tx, reader_done_rx) = mpsc::channel();
std::thread::spawn(move || {
let mut buf = [0u8; 8192];
let mut osc = crate::osc133::Scanner::new();
let mut cmd_started: Option<Instant> = None;
loop {
match reader.read(&mut buf) {
Ok(0) | Err(_) => break,
Ok(n) => {
if let Ok(mut p) = reader_parser.lock() {
p.process(&buf[..n]);
}
if let Some(sess) = &ledger_session {
for ev in osc.feed(&buf[..n]) {
match ev {
crate::osc133::CmdEvent::Start => cmd_started = Some(Instant::now()),
crate::osc133::CmdEvent::End { command, cwd, exit } => {
let dur = cmd_started.take().map(|t| t.elapsed().as_secs());
if let Some(entry) = crate::osc133::to_ledger_entry(
sess, &ledger_surface, command, cwd, exit, dur,
) {
crate::worklog::record(&entry);
}
}
}
}
}
if output_tx.send(TermEvent::Output(id)).is_err() {
break;
}
}
}
}
let _ = reader_done_tx.send(());
});
let exit_code = Arc::new(Mutex::new(None));
let wait_exit_code = exit_code.clone();
let notify_exit = Arc::new(AtomicBool::new(true));
let wait_notify_exit = notify_exit.clone();
let (kill_tx, kill_rx) = mpsc::channel();
std::thread::spawn(move || {
let code = loop {
match child.try_wait() {
Ok(Some(status)) => break Some(status.exit_code() as i32),
Err(_) => break None,
Ok(None) => {}
}
match kill_rx.recv_timeout(Duration::from_millis(20)) {
Ok(()) => {
let _ = child.kill();
break child.wait().ok().map(|status| status.exit_code() as i32);
}
Err(mpsc::RecvTimeoutError::Disconnected) => {
let _ = child.kill();
break child.wait().ok().map(|status| status.exit_code() as i32);
}
Err(mpsc::RecvTimeoutError::Timeout) => {}
}
};
if let Ok(mut slot) = wait_exit_code.lock() {
*slot = code;
}
if wait_notify_exit.swap(false, Ordering::AcqRel) {
let _ = reader_done_rx.recv_timeout(Duration::from_millis(100));
let _ = tx.send(TermEvent::Exited(id));
}
});
let marker = format!(
"__MARS_READY_{:x}__",
NEXT_STARTUP_PROBE.fetch_add(1, Ordering::Relaxed)
);
Ok(Term {
exited: false,
spawn_cwd,
parser,
writer,
master: pair.master,
kill_tx,
startup_input: Some(StartupInput {
bytes: Vec::new(),
marker,
probe_interval: startup_probe_interval,
last_probe: None,
}),
exit_code,
notify_exit,
rows,
cols,
view_offset: 0,
scrollback_limit: scrollback,
})
}
impl Drop for Term {
fn drop(&mut self) {
self.notify_exit.store(false, Ordering::Release);
let _ = self.kill_tx.send(());
}
}
impl Term {
fn prompt_visible(&self) -> bool {
let Ok(parser) = self.parser.lock() else { return false };
let screen = parser.screen();
if screen.hide_cursor() {
return false;
}
let (row, _) = screen.cursor_position();
let (_, cols) = screen.size();
let line: String = (0..cols)
.filter_map(|col| screen.cell(row, col))
.map(|cell| cell.contents())
.collect();
matches!(
line.trim_end().chars().last(),
Some('$' | '#' | '%' | '>' | '❯' | '➜' | 'λ' | '»' | '›')
)
}
pub fn flush_startup_input(&mut self) {
if self.startup_input.is_none() {
return;
}
if self.prompt_visible() {
let bytes = self.startup_input.take().map(|startup| startup.bytes);
if let Some(bytes) = bytes {
self.write_input(&bytes);
}
return;
}
let marker_seen = self.startup_input.as_ref().is_some_and(|startup| {
let Ok(parser) = self.parser.lock() else { return false };
parser
.screen()
.contents()
.lines()
.any(|line| line.trim() == startup.marker)
});
if marker_seen {
let bytes = self.startup_input.take().map(|startup| startup.bytes);
if let Some(bytes) = bytes {
self.write_input(&bytes);
}
return;
}
let probe = self.startup_input.as_mut().and_then(|startup| {
let last_probe = startup.last_probe?;
if last_probe.elapsed() < startup.probe_interval {
return None;
}
startup.last_probe = Some(Instant::now());
Some(format!("echo {}\r", startup.marker))
});
if let Some(probe) = probe {
self.write_input(probe.as_bytes());
}
}
pub fn send_bytes(&mut self, bytes: &[u8]) {
if let Some(startup) = self.startup_input.as_mut() {
startup.bytes.extend_from_slice(bytes);
startup.last_probe.get_or_insert_with(Instant::now);
self.flush_startup_input();
return;
}
self.write_input(bytes);
}
fn write_input(&mut self, bytes: &[u8]) {
let _ = self.writer.write_all(bytes);
let _ = self.writer.flush();
}
pub fn exit_code(&self) -> Option<i32> {
self.exit_code.lock().ok().and_then(|code| *code)
}
pub fn resize(&mut self, rows: u16, cols: u16) {
let rows = rows.max(1);
let cols = cols.max(1);
if rows == self.rows && cols == self.cols {
return;
}
self.rows = rows;
self.cols = cols;
let _ = self.master.resize(PtySize {
rows,
cols,
pixel_width: 0,
pixel_height: 0,
});
if let Ok(mut p) = self.parser.lock() {
p.set_size(rows, cols);
}
}
pub fn screen(&self) -> vt100::Screen {
self.parser.lock().unwrap().screen().clone()
}
pub fn scroll_view(&mut self, delta: i64) {
let requested = (self.view_offset as i64 + delta)
.clamp(0, self.scrollback_limit as i64) as usize;
if let Ok(mut p) = self.parser.lock() {
p.set_scrollback(requested);
self.view_offset = p.screen().scrollback();
} else {
self.view_offset = requested;
}
}
pub fn history_tail(&self, lines: usize) -> String {
let Ok(mut p) = self.parser.lock() else { return String::new() };
let rows = p.screen().size().0 as usize;
let saved = self.view_offset;
let mut pages: Vec<String> = Vec::new();
let (mut off, mut got) = (0usize, 0usize);
loop {
p.set_scrollback(off);
pages.push(p.screen().contents());
got += rows;
if got >= lines || off >= self.scrollback_limit {
break;
}
off += rows;
}
p.set_scrollback(saved); pages.reverse(); let joined = pages.join("\n");
let all: Vec<&str> = joined.lines().collect();
let start = all.len().saturating_sub(lines);
all[start..].join("\n")
}
pub fn scroll_to_live(&mut self) {
if self.view_offset != 0 {
self.scroll_view(-(self.view_offset as i64));
}
}
pub fn view_offset(&self) -> usize {
self.view_offset
}
}