use crossterm::event::{KeyCode, KeyEvent, KeyModifiers};
use std::io::Write;
const SCROLLBACK: usize = 5_000;
const WHEEL_LINES: usize = 3;
pub mod frame {
pub const OUTPUT: u8 = b'o';
pub const SIZE: u8 = b's';
pub const KEYS: u8 = b'k';
pub const RESIZE: u8 = b'r';
const MAX_FRAME: usize = 1 << 20;
const HEADER: usize = 1 + 4;
pub fn encode(kind: u8, payload: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(HEADER + payload.len());
out.push(kind);
out.extend_from_slice(&(payload.len() as u32).to_be_bytes());
out.extend_from_slice(payload);
out
}
pub fn size(kind: u8, cols: u16, rows: u16) -> Vec<u8> {
encode(kind, format!("{cols} {rows}").as_bytes())
}
pub fn parse_size(payload: &[u8]) -> Option<(u16, u16)> {
let text = std::str::from_utf8(payload).ok()?;
let (cols, rows) = text.trim().split_once(' ')?;
let (cols, rows) = (cols.parse().ok()?, rows.parse().ok()?);
(cols > 0 && rows > 0).then_some((cols, rows))
}
#[derive(Default)]
pub struct Decoder {
buf: Vec<u8>,
lost: bool,
}
impl Decoder {
pub fn push(&mut self, bytes: &[u8]) {
self.buf.extend_from_slice(bytes);
}
pub fn next(&mut self) -> Option<(u8, Vec<u8>)> {
if self.lost || self.buf.len() < HEADER {
return None;
}
let len = u32::from_be_bytes(self.buf[1..HEADER].try_into().ok()?) as usize;
if len > MAX_FRAME {
self.lost = true;
return None;
}
if self.buf.len() < HEADER + len {
return None;
}
let kind = self.buf[0];
let payload = self.buf[HEADER..HEADER + len].to_vec();
self.buf.drain(..HEADER + len);
Some((kind, payload))
}
}
}
#[cfg_attr(not(unix), allow(dead_code))]
enum Event {
Output(Vec<u8>),
Size(u16, u16),
}
pub struct Attach {
input: Box<dyn Write + Send>,
close: Box<dyn Fn() + Send>,
pending: std::sync::Arc<std::sync::Mutex<Vec<Event>>>,
closed: std::sync::Arc<std::sync::atomic::AtomicBool>,
pub size: (u16, u16),
requested: (u16, u16),
pub parser: vt100::Parser,
}
impl Attach {
pub fn pump(&mut self) -> bool {
let events = {
let mut pending = self
.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
std::mem::take(&mut *pending)
};
if events.is_empty() {
return false;
}
for event in events {
match event {
Event::Output(bytes) => self.parser.process(&bytes),
Event::Size(cols, rows) => {
self.size = (cols, rows);
self.parser.screen_mut().set_size(rows, cols);
}
}
}
true
}
pub fn closed(&self) -> bool {
self.closed.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn resize(&mut self, cols: u16, rows: u16) -> bool {
if (cols, rows) == self.requested || cols == 0 || rows == 0 {
return true;
}
self.requested = (cols, rows);
self.send(&frame::size(frame::RESIZE, cols, rows))
}
pub fn send_key(&mut self, key: KeyEvent) -> bool {
if self.parser.screen().scrollback() > 0 {
self.parser.screen_mut().set_scrollback(0);
}
match encode(key) {
Some(bytes) => self.send(&frame::encode(frame::KEYS, &bytes)),
None => true,
}
}
pub fn wheel(&mut self, up: bool, col: u16, row: u16) -> bool {
use vt100::{MouseProtocolEncoding as Enc, MouseProtocolMode as Mode};
if self.parser.screen().mouse_protocol_mode() == Mode::None {
let at = self.parser.screen().scrollback();
let to = match up {
true => at + WHEEL_LINES,
false => at.saturating_sub(WHEEL_LINES),
};
self.parser.screen_mut().set_scrollback(to);
return true;
}
let button = 64 + u16::from(!up);
let bytes = match self.parser.screen().mouse_protocol_encoding() {
Enc::Sgr => format!("\x1b[<{button};{};{}M", col + 1, row + 1).into_bytes(),
_ => {
if col > 222 || row > 222 {
return true;
}
let cell = |v: u16| u8::try_from(v + 33).unwrap_or(u8::MAX);
vec![
0x1b,
b'[',
b'M',
u8::try_from(button + 32).unwrap_or(u8::MAX),
cell(col),
cell(row),
]
}
};
self.send(&frame::encode(frame::KEYS, &bytes))
}
fn send(&mut self, bytes: &[u8]) -> bool {
self.input
.write_all(bytes)
.and_then(|()| self.input.flush())
.is_ok()
}
}
impl Drop for Attach {
fn drop(&mut self) {
(self.close)();
}
}
#[cfg(unix)]
pub fn attach(pid: u32) -> Option<Attach> {
use std::io::Read;
let mut stream = connect(pid)?;
let _ = stream.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let mut decoder = frame::Decoder::default();
let mut buf = [0u8; 8192];
let mut early = Vec::new();
let size = loop {
let n = stream.read(&mut buf).ok()?;
if n == 0 {
return None;
}
decoder.push(&buf[..n]);
let mut size = None;
while let Some(event) = read_event(&mut decoder) {
if let Event::Size(cols, rows) = event {
size = Some((cols, rows));
break;
}
early.push(event);
}
if let Some(size) = size {
break size;
}
};
let _ = stream.set_read_timeout(None);
while let Some(event) = read_event(&mut decoder) {
early.push(event);
}
let pending = std::sync::Arc::new(std::sync::Mutex::new(early));
let closed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let collector = std::sync::Arc::clone(&pending);
let hangup = std::sync::Arc::clone(&closed);
let mut reader = stream.try_clone().ok()?;
std::thread::spawn(move || {
let mut buf = [0u8; 8192];
while let Ok(n) = reader.read(&mut buf) {
if n == 0 {
break;
}
decoder.push(&buf[..n]);
let mut pending = collector
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
while let Some(event) = read_event(&mut decoder) {
pending.push(event);
}
}
hangup.store(true, std::sync::atomic::Ordering::Relaxed);
});
let closer = stream.try_clone().ok()?;
Some(Attach {
input: Box::new(stream),
close: Box::new(move || {
let _ = closer.shutdown(std::net::Shutdown::Both);
}),
pending,
closed,
size,
requested: (0, 0),
parser: vt100::Parser::new(size.1, size.0, SCROLLBACK),
})
}
#[cfg(unix)]
fn connect(pid: u32) -> Option<std::os::unix::net::UnixStream> {
let path = crate::shim::socket_path(pid)?;
let mut stream = std::os::unix::net::UnixStream::connect(path).ok()?;
stream.write_all(crate::shim::ATTACH_MAGIC).ok()?;
Some(stream)
}
#[cfg(not(unix))]
pub fn attach(_pid: u32) -> Option<Attach> {
None
}
#[cfg(unix)]
pub fn run_terminal(args: &[String]) -> anyhow::Result<i32> {
let sessions = crate::shim::sessions();
let pid = match args {
[] if sessions.len() == 1 => sessions[0],
[] => return list_sessions(&sessions),
[arg] => arg
.parse::<u32>()
.map_err(|_| anyhow::anyhow!("expected a pid; `cctop attach` lists what is running"))?,
_ => anyhow::bail!("usage: cctop attach [pid]"),
};
if !sessions.contains(&pid) {
anyhow::bail!(
"no cctop-launched agent is running as {pid}; `cctop attach` lists what is running"
);
}
proxy(pid)
}
#[cfg(unix)]
fn list_sessions(sessions: &[u32]) -> anyhow::Result<i32> {
if sessions.is_empty() {
eprintln!(
"No agents are running under cctop. Start one with `cctop claude` \
(or codex, opencode, pi) and it becomes attachable."
);
return Ok(1);
}
eprintln!("Agents running under cctop:\n");
let mut sys = sysinfo::System::new();
let pids: Vec<sysinfo::Pid> = sessions
.iter()
.map(|&p| sysinfo::Pid::from_u32(p))
.collect();
sys.refresh_processes_specifics(
sysinfo::ProcessesToUpdate::Some(&pids),
true,
sysinfo::ProcessRefreshKind::nothing()
.with_cmd(sysinfo::UpdateKind::Always)
.with_cwd(sysinfo::UpdateKind::Always),
);
for &pid in sessions {
let process = sys.process(sysinfo::Pid::from_u32(pid));
let command = process
.map(|p| {
p.cmd()
.iter()
.map(|a| a.to_string_lossy())
.collect::<Vec<_>>()
.join(" ")
})
.unwrap_or_default();
let cwd = process
.and_then(|p| p.cwd())
.map(|d| format!(" in {}", d.display()))
.unwrap_or_default();
eprintln!(" {pid:>7} {}{cwd}", crate::util::truncate(&command, 60));
}
eprintln!("\nAttach with `cctop attach <pid>`.");
Ok(0)
}
#[cfg(unix)]
const DETACH: &[u8] = b"\x1b[24~";
#[cfg(unix)]
const RESTORE: &[u8] = b"\x1b[?1049l\x1b[?1000l\x1b[?1002l\x1b[?1003l\x1b[?1006l\x1b[?25h\x1b[0m";
#[cfg(unix)]
fn proxy(pid: u32) -> anyhow::Result<i32> {
use std::io::{IsTerminal, Read};
if !std::io::stdin().is_terminal() || !std::io::stdout().is_terminal() {
anyhow::bail!("cctop attach needs a terminal; run it from an interactive shell");
}
let stream = connect(pid)
.ok_or_else(|| anyhow::anyhow!("could not open the control socket for {pid}"))?;
let raw = crossterm::terminal::enable_raw_mode().is_ok();
let mut stdout = std::io::stdout();
let _ = stdout.write_all(b"\x1b[?1049h");
let _ = stdout.flush();
let mut input = stream.try_clone()?;
if let Ok((cols, rows)) = crossterm::terminal::size() {
let _ = input.write_all(&frame::size(frame::RESIZE, cols, rows));
}
{
let mut input = input.try_clone()?;
let socket = stream.try_clone()?;
std::thread::spawn(move || {
let mut stdin = std::io::stdin();
let mut buf = [0u8; 1024];
while let Ok(n) = stdin.read(&mut buf) {
if n == 0 || buf[..n].windows(DETACH.len()).any(|w| w == DETACH) {
break;
}
if input
.write_all(&frame::encode(frame::KEYS, &buf[..n]))
.and_then(|()| input.flush())
.is_err()
{
break;
}
}
let _ = socket.shutdown(std::net::Shutdown::Both);
});
}
{
let mut input = input.try_clone()?;
std::thread::spawn(move || {
let mut last = crossterm::terminal::size().unwrap_or((80, 24));
loop {
std::thread::sleep(std::time::Duration::from_millis(500));
let Ok(size) = crossterm::terminal::size() else {
continue;
};
if size != last {
last = size;
if input
.write_all(&frame::size(frame::RESIZE, size.0, size.1))
.is_err()
{
break;
}
}
}
});
}
let mut reader = stream;
let mut decoder = frame::Decoder::default();
let mut buf = [0u8; 8192];
'session: while let Ok(n) = reader.read(&mut buf) {
if n == 0 {
break;
}
decoder.push(&buf[..n]);
while let Some((kind, payload)) = decoder.next() {
if kind == frame::OUTPUT
&& (stdout.write_all(&payload).is_err() || stdout.flush().is_err())
{
break 'session;
}
}
}
let _ = stdout.write_all(RESTORE);
let _ = stdout.flush();
if raw {
let _ = crossterm::terminal::disable_raw_mode();
}
eprintln!("Detached from {pid}. It is still running — `cctop attach {pid}` comes back.");
Ok(0)
}
#[cfg(unix)]
fn read_event(decoder: &mut frame::Decoder) -> Option<Event> {
loop {
let (kind, payload) = decoder.next()?;
match kind {
frame::OUTPUT => return Some(Event::Output(payload)),
frame::SIZE => {
if let Some((cols, rows)) = frame::parse_size(&payload) {
return Some(Event::Size(cols, rows));
}
}
_ => {}
}
}
}
fn encode(key: KeyEvent) -> Option<Vec<u8>> {
let ctrl = key.modifiers.contains(KeyModifiers::CONTROL);
let alt = key.modifiers.contains(KeyModifiers::ALT);
let bytes = match key.code {
KeyCode::Char(c) if ctrl && c.is_ascii_alphabetic() => {
vec![c.to_ascii_lowercase() as u8 & 0x1f]
}
KeyCode::Char(c) => {
let mut b = c.to_string().into_bytes();
if alt {
b.insert(0, 0x1b);
}
b
}
KeyCode::Enter => vec![b'\r'],
KeyCode::Tab => vec![b'\t'],
KeyCode::BackTab => b"\x1b[Z".to_vec(),
KeyCode::Backspace => vec![0x7f],
KeyCode::Esc => vec![0x1b],
KeyCode::Up => b"\x1b[A".to_vec(),
KeyCode::Down => b"\x1b[B".to_vec(),
KeyCode::Right => b"\x1b[C".to_vec(),
KeyCode::Left => b"\x1b[D".to_vec(),
KeyCode::Home => b"\x1b[H".to_vec(),
KeyCode::End => b"\x1b[F".to_vec(),
KeyCode::Insert => b"\x1b[2~".to_vec(),
KeyCode::Delete => b"\x1b[3~".to_vec(),
KeyCode::PageUp => b"\x1b[5~".to_vec(),
KeyCode::PageDown => b"\x1b[6~".to_vec(),
_ => return None,
};
Some(bytes)
}
#[cfg(test)]
mod tests {
use super::*;
fn probe() -> (Attach, std::sync::Arc<std::sync::Mutex<Vec<u8>>>) {
#[derive(Clone)]
struct Sink(std::sync::Arc<std::sync::Mutex<Vec<u8>>>);
impl Write for Sink {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
let sink = Sink(std::sync::Arc::new(std::sync::Mutex::new(Vec::new())));
let written = sink.0.clone();
let attach = Attach {
input: Box::new(sink),
close: Box::new(|| {}),
pending: std::sync::Arc::default(),
closed: std::sync::Arc::default(),
size: (80, 24),
requested: (0, 0),
parser: vt100::Parser::new(24, 80, SCROLLBACK),
};
(attach, written)
}
#[test]
fn the_wheel_reaches_only_an_agent_that_asked_for_it() {
let (mut attach, written) = probe();
assert!(attach.wheel(true, 2, 8));
assert!(written.lock().unwrap().is_empty());
attach.parser.process(b"\x1b[?1000h\x1b[?1006h");
assert!(attach.wheel(true, 2, 8));
assert!(attach.wheel(false, 2, 8));
assert_eq!(
written.lock().unwrap().as_slice(),
[
frame::encode(frame::KEYS, b"\x1b[<64;3;9M"),
frame::encode(frame::KEYS, b"\x1b[<65;3;9M"),
]
.concat()
);
}
#[test]
fn a_wheel_in_the_default_encoding_says_what_it_can() {
let (mut attach, written) = probe();
attach.parser.process(b"\x1b[?1000h");
assert!(attach.wheel(true, 2, 8));
assert_eq!(
written.lock().unwrap().as_slice(),
frame::encode(frame::KEYS, b"\x1b[M\x60\x23\x29")
);
written.lock().unwrap().clear();
assert!(attach.wheel(true, 300, 8));
assert!(written.lock().unwrap().is_empty());
}
#[test]
fn a_pane_with_no_one_reading_the_mouse_scrolls_what_cctop_kept() {
let (mut attach, written) = probe();
for i in 0..100 {
attach.parser.process(format!("line {i}\r\n").as_bytes());
}
assert!(attach.wheel(true, 0, 0));
assert!(attach.wheel(true, 0, 0));
assert_eq!(attach.parser.screen().scrollback(), WHEEL_LINES * 2);
assert!(written.lock().unwrap().is_empty());
assert!(attach.wheel(false, 0, 0));
assert_eq!(attach.parser.screen().scrollback(), WHEEL_LINES);
attach.send_key(KeyEvent::from(KeyCode::Char('y')));
assert_eq!(attach.parser.screen().scrollback(), 0);
assert!(attach.wheel(false, 0, 0));
assert_eq!(attach.parser.screen().scrollback(), 0);
}
#[test]
fn a_size_payload_is_read_or_refused() {
assert_eq!(frame::parse_size(b"120 40"), Some((120, 40)));
assert_eq!(frame::parse_size(b"0 40"), None);
assert_eq!(frame::parse_size(b"something else"), None);
assert_eq!(frame::parse_size(b""), None);
}
#[test]
fn frames_survive_a_stream_that_splits_them_anywhere() {
let mut wire = Vec::new();
wire.extend(frame::size(frame::SIZE, 100, 30));
wire.extend(frame::encode(
frame::OUTPUT,
b"o\x00\x00\x00\x09hello\x1b[A",
));
wire.extend(frame::encode(frame::KEYS, b""));
let mut decoder = frame::Decoder::default();
let mut got = Vec::new();
for byte in &wire {
decoder.push(&[*byte]);
while let Some(f) = decoder.next() {
got.push(f);
}
}
assert_eq!(
got,
vec![
(frame::SIZE, b"100 30".to_vec()),
(frame::OUTPUT, b"o\x00\x00\x00\x09hello\x1b[A".to_vec()),
(frame::KEYS, Vec::new()),
]
);
let mut decoder = frame::Decoder::default();
decoder.push(&wire);
assert_eq!(std::iter::from_fn(|| decoder.next()).count(), 3);
}
#[test]
fn a_desynchronised_stream_yields_nothing_rather_than_garbage() {
let mut decoder = frame::Decoder::default();
decoder.push(&[frame::OUTPUT, 0xff, 0xff, 0xff, 0xff]);
decoder.push(b"whatever follows");
assert_eq!(decoder.next(), None);
decoder.push(&frame::size(frame::SIZE, 80, 24));
assert_eq!(decoder.next(), None);
}
#[test]
fn keys_encode_as_a_terminal_would_send_them() {
let plain = |c| KeyEvent::new(KeyCode::Char(c), KeyModifiers::NONE);
let ctrl = |c| KeyEvent::new(KeyCode::Char(c), KeyModifiers::CONTROL);
let key = |code| KeyEvent::new(code, KeyModifiers::NONE);
assert_eq!(encode(plain('a')), Some(b"a".to_vec()));
assert_eq!(encode(plain('é')), Some("é".as_bytes().to_vec()));
assert_eq!(encode(ctrl('c')), Some(vec![3]));
assert_eq!(encode(ctrl('C')), Some(vec![3]));
assert_eq!(encode(key(KeyCode::Enter)), Some(b"\r".to_vec()));
assert_eq!(encode(key(KeyCode::Backspace)), Some(vec![0x7f]));
assert_eq!(encode(key(KeyCode::Up)), Some(b"\x1b[A".to_vec()));
assert_eq!(encode(key(KeyCode::F(5))), None);
}
}