use std::fs::File;
use std::io::{Read, Write};
use std::os::fd::{AsRawFd, FromRawFd};
use std::os::unix::net::UnixStream;
use std::os::unix::process::CommandExt;
use std::path::PathBuf;
use std::process::Command;
use std::sync::mpsc::SyncSender;
use std::sync::{Arc, Mutex};
const RESIZE_POLL: std::time::Duration = std::time::Duration::from_millis(500);
pub fn run(argv: &[String]) -> anyhow::Result<i32> {
if argv.is_empty() {
anyhow::bail!("usage: cctop run <command> [args…] (e.g. cctop run claude)");
}
let (mut child, master) = spawn_on_pty(argv, None)?;
let pid = child.id();
let listener = listen(pid)?;
let socket = socket_path(pid);
let raw = crossterm::terminal::enable_raw_mode().is_ok();
let has_window = std::io::IsTerminal::is_terminal(&std::io::stdout());
let local = match has_window {
true => crossterm::terminal::size().unwrap_or((80, 24)),
false => (0, 0),
};
let fan = Arc::new(Mutex::new(Fanout::new(master.try_clone()?, local)));
{
let (master, fan) = (master.try_clone()?, Arc::clone(&fan));
std::thread::spawn(move || pump_output(master, fan, Echo::Yes));
}
{
let master = master.try_clone()?;
std::thread::spawn(move || pump_input(master));
}
if has_window {
let fan = Arc::clone(&fan);
std::thread::spawn(move || watch_resize(fan));
}
let control = master.try_clone()?;
std::thread::spawn(move || serve(listener, control, fan));
let status = child.wait();
if raw {
let _ = crossterm::terminal::disable_raw_mode();
}
if let Some(socket) = socket {
let _ = std::fs::remove_file(socket);
}
Ok(status?.code().unwrap_or(1))
}
pub struct Hosted {
pub pid: u32,
pub label: String,
child: std::process::Child,
socket: Option<PathBuf>,
}
impl Hosted {
pub fn finished(&mut self) -> Option<i32> {
match self.child.try_wait() {
Ok(Some(status)) => Some(status.code().unwrap_or(1)),
Err(_) => Some(1),
Ok(None) => None,
}
}
}
impl Drop for Hosted {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
if let Some(socket) = &self.socket {
let _ = std::fs::remove_file(socket);
}
}
}
pub fn host(argv: &[String], cwd: Option<&std::path::Path>) -> anyhow::Result<Hosted> {
if argv.is_empty() {
anyhow::bail!("usage: cctop <command> [args…] (e.g. cctop claude)");
}
let (child, master) = spawn_on_pty(argv, cwd)?;
let pid = child.id();
let listener = listen(pid)?;
let socket = socket_path(pid);
let fan = Arc::new(Mutex::new(Fanout::new(master.try_clone()?, (0, 0))));
{
let (master, fan) = (master.try_clone()?, Arc::clone(&fan));
std::thread::spawn(move || pump_output(master, fan, Echo::No));
}
std::thread::spawn(move || serve(listener, master, fan));
let label = argv
.iter()
.map(|arg| arg.rsplit('/').next().unwrap_or(arg))
.collect::<Vec<_>>()
.join(" ");
Ok(Hosted {
pid,
label,
child,
socket,
})
}
pub fn is_command(word: &str) -> bool {
use std::os::unix::fs::PermissionsExt;
let executable = |p: PathBuf| {
std::fs::metadata(&p).is_ok_and(|m| m.is_file() && m.permissions().mode() & 0o111 != 0)
};
if word.contains('/') {
return executable(PathBuf::from(word));
}
std::env::split_paths(&std::env::var_os("PATH").unwrap_or_default())
.any(|dir| executable(dir.join(word)))
}
pub fn socket_path(pid: u32) -> Option<PathBuf> {
Some(base_dir()?.join(format!("{pid}.sock")))
}
fn base_dir() -> Option<PathBuf> {
dirs::runtime_dir()
.or_else(dirs::cache_dir)
.map(|d| d.join("cctop"))
}
pub fn sessions() -> Vec<u32> {
let Some(dir) = base_dir() else {
return Vec::new();
};
let Ok(entries) = std::fs::read_dir(dir) else {
return Vec::new();
};
let mut found: Vec<(std::time::SystemTime, u32)> = Vec::new();
for entry in entries.flatten() {
let path = entry.path();
let Some(pid) = path
.file_name()
.and_then(|n| n.to_str())
.and_then(|n| n.strip_suffix(".sock"))
.and_then(|n| n.parse::<u32>().ok())
else {
continue;
};
if UnixStream::connect(&path).is_err() {
let _ = std::fs::remove_file(&path);
continue;
}
let created = entry
.metadata()
.and_then(|m| m.created().or_else(|_| m.modified()))
.unwrap_or(std::time::UNIX_EPOCH);
found.push((created, pid));
}
found.sort_unstable();
found.into_iter().map(|(_, pid)| pid).collect()
}
fn listen(pid: u32) -> anyhow::Result<std::os::unix::net::UnixListener> {
let path = socket_path(pid).ok_or_else(|| anyhow::anyhow!("no runtime or cache directory"))?;
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
let _ =
std::fs::set_permissions(parent, std::os::unix::fs::PermissionsExt::from_mode(0o700));
}
let _ = std::fs::remove_file(&path);
Ok(std::os::unix::net::UnixListener::bind(&path)?)
}
pub const ATTACH_MAGIC: &[u8] = b"\x00cctop-attach\n";
const REPLAY_BYTES: usize = 256 * 1024;
struct Sub {
id: u64,
tx: SyncSender<Vec<u8>>,
size: (u16, u16),
}
struct Fanout {
master: File,
recent: Vec<u8>,
subs: Vec<Sub>,
next_id: u64,
local: (u16, u16),
applied: (u16, u16),
}
impl Fanout {
fn new(master: File, local: (u16, u16)) -> Self {
Self {
recent: Vec::new(),
subs: Vec::new(),
next_id: 0,
local,
applied: pty_size(&master),
master,
}
}
fn push(&mut self, chunk: &[u8]) {
self.recent.extend_from_slice(chunk);
if self.recent.len() > REPLAY_BYTES {
self.recent.drain(..self.recent.len() - REPLAY_BYTES);
}
let frame = crate::attach::frame::encode(crate::attach::frame::OUTPUT, chunk);
let before = self.subs.len();
self.subs
.retain(|sub| sub.tx.try_send(frame.clone()).is_ok());
if self.subs.len() != before {
self.fit();
}
}
fn set_local(&mut self, size: (u16, u16)) {
self.local = size;
self.fit();
}
fn remove(&mut self, id: u64) {
self.subs.retain(|sub| sub.id != id);
self.fit();
}
fn set_sub_size(&mut self, id: u64, size: (u16, u16)) {
if let Some(sub) = self.subs.iter_mut().find(|sub| sub.id == id) {
sub.size = size;
}
self.fit();
}
fn fit(&mut self) {
let size = std::iter::once(self.local)
.chain(self.subs.iter().map(|sub| sub.size))
.filter(|&(cols, rows)| cols > 0 && rows > 0)
.fold(None::<(u16, u16)>, |acc, (cols, rows)| match acc {
Some((c, r)) => Some((c.min(cols), r.min(rows))),
None => Some((cols, rows)),
});
let Some((cols, rows)) = size else { return };
if (cols, rows) == self.applied {
return;
}
self.applied = (cols, rows);
let ws = winsize(cols, rows);
unsafe { libc::ioctl(self.master.as_raw_fd(), libc::TIOCSWINSZ as _, &ws) };
let frame = crate::attach::frame::size(crate::attach::frame::SIZE, cols, rows);
self.subs
.retain(|sub| sub.tx.try_send(frame.clone()).is_ok());
}
}
fn locked(fan: &Mutex<Fanout>) -> std::sync::MutexGuard<'_, Fanout> {
fan.lock().unwrap_or_else(|e| e.into_inner())
}
fn serve(listener: std::os::unix::net::UnixListener, master: File, fan: Arc<Mutex<Fanout>>) {
for stream in listener.incoming().flatten() {
let Ok(master) = master.try_clone() else {
continue;
};
let fan = Arc::clone(&fan);
std::thread::spawn(move || converse(stream, master, fan));
}
}
fn converse(mut stream: UnixStream, mut master: File, fan: Arc<Mutex<Fanout>>) {
use crate::attach::frame;
let mut buf = [0u8; 4096];
let mut first = true;
let mut watcher: Option<(u64, frame::Decoder)> = None;
'connection: while let Ok(n) = stream.read(&mut buf) {
if n == 0 {
break;
}
let mut bytes = &buf[..n];
if std::mem::take(&mut first)
&& let Some(rest) = bytes.strip_prefix(ATTACH_MAGIC)
{
let Some(id) = subscribe(&stream, &fan) else {
break;
};
watcher = Some((id, frame::Decoder::default()));
bytes = rest;
}
let Some((id, decoder)) = watcher.as_mut() else {
if !bytes.is_empty() && (master.write_all(bytes).is_err() || master.flush().is_err()) {
break;
}
continue;
};
decoder.push(bytes);
while let Some((kind, payload)) = decoder.next() {
let delivered = match kind {
frame::KEYS => master
.write_all(&payload)
.and_then(|()| master.flush())
.is_ok(),
frame::RESIZE => {
if let Some(size) = frame::parse_size(&payload) {
locked(&fan).set_sub_size(*id, size);
}
true
}
_ => true,
};
if !delivered {
break 'connection;
}
}
}
if let Some((id, _)) = watcher {
locked(&fan).remove(id);
}
}
fn subscribe(stream: &UnixStream, fan: &Arc<Mutex<Fanout>>) -> Option<u64> {
use crate::attach::frame;
let mut out = stream.try_clone().ok()?;
let (tx, rx) = std::sync::mpsc::sync_channel::<Vec<u8>>(256);
let id = {
let mut fan = locked(fan);
let (cols, rows) = fan.applied;
let header = frame::size(frame::SIZE, cols, rows);
let replay = frame::encode(frame::OUTPUT, &fan.recent);
if out.write_all(&header).is_err() || out.write_all(&replay).is_err() {
return None;
}
let id = fan.next_id;
fan.next_id += 1;
fan.subs.push(Sub {
id,
tx,
size: (0, 0),
});
id
};
std::thread::spawn(move || {
while let Ok(chunk) = rx.recv() {
if out.write_all(&chunk).is_err() {
break;
}
}
});
Some(id)
}
fn pty_size(master: &File) -> (u16, u16) {
let mut ws = winsize(80, 24);
if unsafe { libc::ioctl(master.as_raw_fd(), libc::TIOCGWINSZ as _, &raw mut ws) } == -1 {
return (80, 24);
}
(ws.ws_col, ws.ws_row)
}
fn spawn_on_pty(
argv: &[String],
cwd: Option<&std::path::Path>,
) -> anyhow::Result<(std::process::Child, File)> {
let (cols, rows) = crossterm::terminal::size().unwrap_or((80, 24));
let mut size = winsize(cols, rows);
let mut master_fd = -1;
let mut slave_fd = -1;
let rc = unsafe {
libc::openpty(
&mut master_fd,
&mut slave_fd,
std::ptr::null_mut(),
std::ptr::null_mut::<libc::termios>(),
&raw mut size,
)
};
if rc != 0 {
return Err(std::io::Error::last_os_error().into());
}
let master = unsafe { File::from_raw_fd(master_fd) };
let slave = unsafe { File::from_raw_fd(slave_fd) };
let mut cmd = Command::new(&argv[0]);
cmd.args(&argv[1..]);
cmd.env_remove("TMUX");
cmd.env_remove("TMUX_PANE");
if let Some(cwd) = cwd.filter(|dir| dir.is_dir()) {
cmd.current_dir(cwd);
}
unsafe {
cmd.pre_exec(move || {
if libc::setsid() == -1 || libc::ioctl(slave_fd, libc::TIOCSCTTY as _, 0) == -1 {
return Err(std::io::Error::last_os_error());
}
for target in 0..3 {
if libc::dup2(slave_fd, target) == -1 {
return Err(std::io::Error::last_os_error());
}
}
if slave_fd > 2 {
libc::close(slave_fd);
}
libc::close(master_fd);
Ok(())
})
};
let child = cmd.spawn()?;
drop(slave);
Ok((child, master))
}
#[cfg(all(test, target_os = "linux"))]
pub(crate) fn test_session(argv: &[&str], local: (u16, u16)) -> (std::process::Child, u32) {
let argv: Vec<String> = argv.iter().map(|s| (*s).to_string()).collect();
let (child, master) = spawn_on_pty(&argv, None).expect("pty child");
let pid = child.id();
let listener = listen(pid).expect("control socket");
let fan = Arc::new(Mutex::new(Fanout::new(
master.try_clone().expect("master"),
local,
)));
locked(&fan).fit();
{
let (master, fan) = (master.try_clone().expect("master"), Arc::clone(&fan));
std::thread::spawn(move || pump_output(master, fan, Echo::No));
}
std::thread::spawn(move || serve(listener, master, fan));
(child, pid)
}
#[derive(PartialEq)]
enum Echo {
Yes,
No,
}
fn pump_output(mut master: File, fan: Arc<Mutex<Fanout>>, echo: Echo) {
let mut out = std::io::stdout();
let mut buf = [0u8; 8192];
let mut answers = Answers::default();
let mut reply = master.try_clone().ok();
while let Ok(n) = master.read(&mut buf) {
if n == 0 {
break;
}
if echo == Echo::Yes && (out.write_all(&buf[..n]).is_err() || out.flush().is_err()) {
break;
}
if echo == Echo::No
&& let Some(reply) = reply.as_mut()
{
answers.answer(&buf[..n], reply);
}
locked(&fan).push(&buf[..n]);
}
}
#[derive(Default)]
struct Answers {
tail: Vec<u8>,
}
const QUERY_MAX: usize = 16;
impl Answers {
fn answer(&mut self, chunk: &[u8], reply: &mut impl Write) {
let mut seen = std::mem::take(&mut self.tail);
seen.extend_from_slice(chunk);
for (query, answer) in Self::table() {
for _ in 0..strike(&mut seen, query) {
let _ = reply.write_all(answer.as_bytes());
}
}
let _ = reply.flush();
let keep = seen.len().saturating_sub(QUERY_MAX);
self.tail = seen.split_off(keep);
}
fn table() -> Vec<(&'static [u8], String)> {
let da1 = "\x1b[?62;22c".to_string();
vec![
(b"\x1b[c".as_slice(), da1.clone()),
(b"\x1b[0c".as_slice(), da1),
(b"\x1b[>c".as_slice(), "\x1b[>0;0;0c".to_string()),
(b"\x1b[>0c".as_slice(), "\x1b[>0;0;0c".to_string()),
(b"\x1b[5n".as_slice(), "\x1b[0n".to_string()),
(b"\x1b[?2026$p".as_slice(), "\x1b[?2026;2$y".to_string()),
(b"\x1b]10;?".as_slice(), color_reply(10)),
(b"\x1b]11;?".as_slice(), color_reply(11)),
]
}
}
fn strike(haystack: &mut [u8], needle: &[u8]) -> usize {
if needle.is_empty() || haystack.len() < needle.len() {
return 0;
}
let mut found = 0;
let mut i = 0;
while i + needle.len() <= haystack.len() {
if &haystack[i..i + needle.len()] == needle {
haystack[i..i + needle.len()].fill(0);
i += needle.len();
found += 1;
continue;
}
i += 1;
}
found
}
fn color_reply(which: u16) -> String {
let light = crate::ui::theme::variant() == crate::ui::theme::Variant::Light;
let bright = "ffff/ffff/ffff";
let dark = "0000/0000/0000";
let value = match (which, light) {
(11, true) | (10, false) => bright,
_ => dark,
};
format!("\x1b]{which};rgb:{value}\x1b\\")
}
fn pump_input(mut master: File) {
let mut stdin = std::io::stdin();
let mut buf = [0u8; 1024];
while let Ok(n) = stdin.read(&mut buf) {
if n == 0 || master.write_all(&buf[..n]).is_err() {
break;
}
}
}
fn watch_resize(fan: Arc<Mutex<Fanout>>) {
let mut last = crossterm::terminal::size().unwrap_or((80, 24));
loop {
std::thread::sleep(RESIZE_POLL);
let Ok(size) = crossterm::terminal::size() else {
continue;
};
if size != last {
last = size;
locked(&fan).set_local(size);
}
}
}
fn winsize(cols: u16, rows: u16) -> libc::winsize {
libc::winsize {
ws_row: rows,
ws_col: cols,
ws_xpixel: 0,
ws_ypixel: 0,
}
}
#[cfg(all(test, target_os = "linux"))]
mod tests {
use super::*;
#[test]
fn a_command_is_told_apart_from_a_stray_word() {
assert!(is_command("sh"));
assert!(is_command("/bin/sh"));
assert!(!is_command("cctop-no-such-command"));
assert!(!is_command("/etc/hostname"));
}
#[test]
fn a_line_sent_to_the_socket_becomes_the_childs_input() {
let out = std::env::temp_dir().join("cctop-shim-test.txt");
let _ = std::fs::remove_file(&out);
let (mut child, master) = spawn_on_pty(
&[
"sh".into(),
"-c".into(),
format!("tee {} >/dev/null; :", out.display()),
],
None,
)
.unwrap();
let pid = child.id();
let listener = listen(pid).unwrap();
let fan = Arc::new(Mutex::new(Fanout::new(
master.try_clone().unwrap(),
(80, 24),
)));
std::thread::spawn(move || serve(listener, master, fan));
let sent = crate::inject::send_line(pid, "continue");
let text = (0..50).find_map(|_| {
std::thread::sleep(std::time::Duration::from_millis(100));
std::fs::read_to_string(&out).ok().filter(|t| !t.is_empty())
});
let _ = child.kill();
let _ = child.wait();
let _ = std::fs::remove_file(&out);
let _ = socket_path(pid).map(std::fs::remove_file);
sent.unwrap();
assert_eq!(text.unwrap().trim_end(), "continue");
}
#[test]
fn a_watcher_gets_the_screen_and_can_still_type() {
use crossterm::event::{KeyCode, KeyEvent, KeyModifiers};
let (mut child, master) = spawn_on_pty(
&[
"sh".into(),
"-c".into(),
"printf 'ALREADY-DRAWN\\r\\n'; cat".into(),
],
None,
)
.unwrap();
let pid = child.id();
let listener = listen(pid).unwrap();
let fan = Arc::new(Mutex::new(Fanout::new(
master.try_clone().unwrap(),
(80, 24),
)));
for spawn in [0, 1] {
let (master, fan) = (master.try_clone().unwrap(), Arc::clone(&fan));
match spawn {
0 => std::thread::spawn(move || pump_output(master, fan, Echo::No)),
_ => {
let listener = listener.try_clone().unwrap();
std::thread::spawn(move || serve(listener, master, fan))
}
};
}
std::thread::sleep(std::time::Duration::from_millis(300));
let mut attach = crate::attach::attach(pid).expect("no attach connection");
let screen = |attach: &mut crate::attach::Attach, want: &str| {
(0..50).any(|_| {
std::thread::sleep(std::time::Duration::from_millis(100));
attach.pump();
attach.parser.screen().contents().contains(want)
})
};
let replayed = screen(&mut attach, "ALREADY-DRAWN");
attach.send_key(KeyEvent::new(KeyCode::Char('z'), KeyModifiers::NONE));
attach.send_key(KeyEvent::new(KeyCode::Enter, KeyModifiers::NONE));
let typed = screen(&mut attach, "z");
let _ = child.kill();
let _ = child.wait();
let _ = socket_path(pid).map(std::fs::remove_file);
assert!(
replayed,
"the replay did not carry what the child had drawn"
);
assert!(typed, "the keystroke never reached the child");
}
#[test]
fn a_watchers_size_reaches_the_agent_and_is_returned_when_it_leaves() {
let (mut child, master) = spawn_on_pty(
&[
"sh".into(),
"-c".into(),
"while :; do stty size; sleep 0.2; done".into(),
],
None,
)
.unwrap();
let pid = child.id();
let listener = listen(pid).unwrap();
let local = (100u16, 40u16);
let fan = Arc::new(Mutex::new(Fanout::new(master.try_clone().unwrap(), local)));
locked(&fan).fit();
{
let (master, fan) = (master.try_clone().unwrap(), Arc::clone(&fan));
std::thread::spawn(move || pump_output(master, fan, Echo::No));
}
{
let (master, fan) = (master.try_clone().unwrap(), Arc::clone(&fan));
std::thread::spawn(move || serve(listener, master, fan));
}
let settled = |want: (u16, u16)| {
(0..50).any(|_| {
std::thread::sleep(std::time::Duration::from_millis(100));
pty_size(&master) == want
})
};
let started_at_local = settled(local);
let mut attach = crate::attach::attach(pid).expect("no attach connection");
attach.resize(60, 20);
let shrank = settled((60, 20));
let agent_told = (0..50).any(|_| {
std::thread::sleep(std::time::Duration::from_millis(100));
attach.pump();
attach.parser.screen().contents().contains("20 60")
});
let watcher_told = attach.size == (60, 20);
drop(attach);
let restored = settled(local);
let _ = child.kill();
let _ = child.wait();
let _ = socket_path(pid).map(std::fs::remove_file);
assert!(started_at_local, "the pty did not start at the local size");
assert!(shrank, "the watcher's size never reached the pty");
assert!(agent_told, "the agent was not told its new size");
assert!(watcher_told, "the watcher was not told the granted size");
assert!(restored, "detaching did not give the window its size back");
}
#[test]
fn a_pty_child_without_a_socket_reports_the_missing_precondition() {
if unsafe { libc::geteuid() } == 0 {
eprintln!("skipping: running as root, so no precondition is missing");
return;
}
let (mut child, _master) =
spawn_on_pty(&["sh".into(), "-c".into(), "sleep 30".into()], None).expect("pty child");
let error = crate::inject::send_line(child.id(), "continue").unwrap_err();
let _ = child.kill();
let _ = child.wait();
assert!(
error.contains("legacy_tiocsti") || error.contains("root"),
"expected a named precondition, got: {error}"
);
}
#[test]
fn a_question_split_across_two_reads_is_answered_exactly_once() {
let mut answers = Answers::default();
let mut said = Vec::new();
answers.answer(b"drawing\x1b", &mut said);
assert!(said.is_empty(), "answered half a question");
answers.answer(b"[c and on", &mut said);
assert_eq!(said, b"\x1b[?62;22c");
let mut answers = Answers::default();
let mut said = Vec::new();
answers.answer(b"\x1b[5n", &mut said);
answers.answer(b"more output", &mut said);
answers.answer(b"and more", &mut said);
assert_eq!(said, b"\x1b[0n", "the tail was answered twice");
}
#[test]
fn the_questions_a_harness_asks_at_startup_get_answers() {
let mut answers = Answers::default();
let mut said = Vec::new();
answers.answer(b"\x1b[?2004h\x1b[>4;0m\x1b[>7u\x1b[?1004h\x1b[6n\x1b]10;?\x1b\\\x1b]11;?\x1b\\\x1b[?u\x1b[c", &mut said);
let said = String::from_utf8_lossy(&said).to_string();
assert!(said.contains("[?62;22c"), "no device attributes: {said:?}");
assert!(said.contains("]10;rgb:"), "no foreground colour: {said:?}");
assert!(said.contains("]11;rgb:"), "no background colour: {said:?}");
assert!(
!said.contains('R'),
"answered the cursor position: {said:?}"
);
}
#[test]
fn a_hosted_agent_gets_an_answer_from_the_shim() {
let script = "stty raw -echo; printf '\\033[c'; dd bs=1 count=9 2>/dev/null | tr -d '\\033'; sleep 30";
let argv: Vec<String> = ["sh", "-c", script].iter().map(|s| s.to_string()).collect();
let Ok(hosted) = host(&argv, None) else {
eprintln!("skipping: no pty available");
return;
};
let mut view = crate::attach::attach(hosted.pid).expect("attach to the shim");
let seen = (0..50).find_map(|_| {
view.pump();
let screen = view.parser.screen().contents();
if screen.contains("[?62;22c") {
return Some(screen);
}
std::thread::sleep(std::time::Duration::from_millis(100));
None
});
drop(hosted);
assert!(
seen.is_some(),
"the agent's question went unanswered on a real pty"
);
}
}