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];
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;
}
locked(&fan).push(&buf[..n]);
}
}
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}"
);
}
}