use std::io::{BufRead, BufReader, Write};
use std::sync::mpsc::{Receiver, TryRecvError, channel};
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
use portable_pty::{Child, CommandBuilder, MasterPty, PtySize, native_pty_system};
use super::Tmux;
use super::control::{Decoder, Event};
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error("could not open a pty: {0}")]
Pty(String),
#[error("could not start tmux: {0}")]
Spawn(String),
#[error("write to tmux failed: {0}")]
Write(#[source] std::io::Error),
#[error("the control connection has closed")]
Closed,
#[error("timed out after {0:?} waiting for tmux")]
Timeout(Duration),
}
pub type Result<T> = std::result::Result<T, Error>;
pub struct ControlClient {
writer: Box<dyn Write + Send>,
events: Receiver<Event>,
child: Box<dyn Child + Send + Sync>,
master: Box<dyn MasterPty + Send>,
reader: Option<JoinHandle<()>>,
}
impl ControlClient {
pub fn attach(tmux: &Tmux, session: &str, size: (u16, u16)) -> Result<Self> {
let (cols, rows) = size;
let pair = native_pty_system()
.openpty(PtySize {
rows,
cols,
pixel_width: 0,
pixel_height: 0,
})
.map_err(|e| Error::Pty(e.to_string()))?;
let mut cmd = CommandBuilder::new("tmux");
for arg in tmux.socket().args() {
cmd.arg(arg);
}
cmd.arg("-CC");
cmd.arg("attach");
cmd.arg("-t");
cmd.arg(format!("={session}"));
cmd.env("TERM", "xterm-256color");
let child = pair
.slave
.spawn_command(cmd)
.map_err(|e| Error::Spawn(e.to_string()))?;
drop(pair.slave);
let read_half = pair
.master
.try_clone_reader()
.map_err(|e| Error::Pty(e.to_string()))?;
let writer = pair
.master
.take_writer()
.map_err(|e| Error::Pty(e.to_string()))?;
let (tx, events) = channel();
let reader = std::thread::spawn(move || {
let mut decoder = Decoder::new();
let mut buf = BufReader::new(read_half);
let mut line = Vec::new();
loop {
line.clear();
match buf.read_until(b'\n', &mut line) {
Ok(0) => break,
Ok(_) => {}
Err(_) => break,
}
if let Some(event) = decoder.push(&line)
&& tx.send(event).is_err()
{
break;
}
}
});
let mut client = Self {
writer,
events,
child,
master: pair.master,
reader: Some(reader),
};
client.resize(cols, rows)?;
Ok(client)
}
pub fn send_command(&mut self, command: &str) -> Result<()> {
self.writer
.write_all(command.as_bytes())
.and_then(|()| self.writer.write_all(b"\n"))
.and_then(|()| self.writer.flush())
.map_err(Error::Write)
}
pub fn try_event(&self) -> std::result::Result<Event, TryRecvError> {
self.events.try_recv()
}
pub fn next_event(&self, timeout: Duration) -> Result<Event> {
self.events
.recv_timeout(timeout)
.map_err(|_| Error::Timeout(timeout))
}
pub fn wait_for(
&self,
timeout: Duration,
mut predicate: impl FnMut(&Event) -> bool,
) -> Result<Vec<Event>> {
let deadline = Instant::now() + timeout;
let mut seen = Vec::new();
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Err(Error::Timeout(timeout));
}
match self.events.recv_timeout(remaining) {
Ok(event) => {
let matched = predicate(&event);
seen.push(event);
if matched {
return Ok(seen);
}
}
Err(_) => return Err(Error::Timeout(timeout)),
}
}
}
pub fn resize(&mut self, cols: u16, rows: u16) -> Result<()> {
self.master
.resize(PtySize {
rows,
cols,
pixel_width: 0,
pixel_height: 0,
})
.map_err(|e| Error::Pty(e.to_string()))?;
self.send_command(&format!("refresh-client -C {cols}x{rows}"))
}
pub fn shutdown(&mut self) -> Result<()> {
let _ = self.send_command("detach-client");
let _ = self.child.kill();
let _ = self.child.wait();
if let Some(handle) = self.reader.take() {
let _ = handle.join();
}
Ok(())
}
}
impl Drop for ControlClient {
fn drop(&mut self) {
let _ = self.shutdown();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::tmux::testing::TestServer;
use crate::tmux::{DEFAULT_SIZE, control};
use tempfile::TempDir;
const TIMEOUT: Duration = Duration::from_secs(10);
fn session(server: &TestServer, name: &str) -> TempDir {
let tmp = TempDir::new().unwrap();
server
.tmux
.new_session(name, tmp.path(), DEFAULT_SIZE)
.unwrap();
tmp
}
#[test]
fn attaches_and_reports_the_session() {
let server = TestServer::new();
let _dir = session(&server, "ctl");
let client = ControlClient::attach(&server.tmux, "ctl", DEFAULT_SIZE).unwrap();
let seen = client
.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::SessionChanged { .. })
})
.unwrap();
let Some(control::Event::SessionChanged { name, .. }) = seen.last() else {
panic!("expected a session-changed event, saw {seen:?}");
};
assert_eq!(name, "ctl");
}
fn pane_size(tmux: &Tmux, session: &str, want: &str) -> String {
let deadline = Instant::now() + TIMEOUT;
let mut last = String::new();
while Instant::now() < deadline {
last = tmux
.run(&[
"list-panes",
"-t",
&format!("={session}"),
"-F",
"#{pane_width}x#{pane_height}",
])
.unwrap_or_default()
.trim()
.to_string();
if last == want {
break;
}
std::thread::sleep(Duration::from_millis(50));
}
last
}
#[test]
fn resizing_moves_the_real_tmux_pane_not_just_the_pty() {
let server = TestServer::new();
let _dir = session(&server, "rz");
let mut client = ControlClient::attach(&server.tmux, "rz", DEFAULT_SIZE).unwrap();
client
.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::SessionChanged { .. })
})
.unwrap();
client.resize(80, 24).unwrap();
assert_eq!(
pane_size(&server.tmux, "rz", "80x24"),
"80x24",
"the pane must follow the view, or the agent's output is garbled"
);
}
#[test]
fn attaching_sizes_the_session_to_the_client() {
let server = TestServer::new();
let _dir = session(&server, "at");
let client = ControlClient::attach(&server.tmux, "at", (100, 30)).unwrap();
client
.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::SessionChanged { .. })
})
.unwrap();
assert_eq!(pane_size(&server.tmux, "at", "100x30"), "100x30");
}
#[test]
fn a_command_reply_comes_back_intact() {
let server = TestServer::new();
let _dir = session(&server, "cmd");
let mut client = ControlClient::attach(&server.tmux, "cmd", DEFAULT_SIZE).unwrap();
client.send_command("list-panes -F \"#{pane_id}\"").unwrap();
let seen = client
.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::CommandReply { lines, .. }
if lines.iter().any(|l| l.starts_with('%')))
})
.unwrap();
let Some(control::Event::CommandReply { lines, error, .. }) = seen.last() else {
panic!("expected a reply, saw {seen:?}");
};
assert!(!error);
assert!(lines[0].starts_with('%'), "got {lines:?}");
}
#[test]
fn a_failing_command_is_reported_as_an_error_reply() {
let server = TestServer::new();
let _dir = session(&server, "bad");
let mut client = ControlClient::attach(&server.tmux, "bad", DEFAULT_SIZE).unwrap();
client.send_command("select-window -t @999").unwrap();
let seen = client
.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::CommandReply { error: true, .. })
})
.unwrap();
assert!(matches!(
seen.last(),
Some(control::Event::CommandReply { error: true, .. })
));
}
#[test]
fn pane_output_is_delivered_as_bytes() {
let server = TestServer::new();
let _dir = session(&server, "out");
let mut client = ControlClient::attach(&server.tmux, "out", DEFAULT_SIZE).unwrap();
let _ = client.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::SessionChanged { .. })
});
let pane = server.tmux.list_panes("out").unwrap().remove(0);
client
.send_command(&format!("send-keys -t {pane} -l 'printf MARVEROUT'"))
.unwrap();
client
.send_command(&format!("send-keys -t {pane} Enter"))
.unwrap();
let mut seen_bytes: Vec<u8> = Vec::new();
let result = client.wait_for(TIMEOUT, |event| {
if let control::Event::Output { data, .. } = event {
seen_bytes.extend_from_slice(data);
}
String::from_utf8_lossy(&seen_bytes).contains("MARVEROUT")
});
assert!(
result.is_ok(),
"never saw the marker; got {:?}",
String::from_utf8_lossy(&seen_bytes)
);
}
#[test]
fn non_ascii_pane_output_arrives_intact_over_the_real_transport() {
const HEARTS: usize = 4000;
let server = TestServer::new();
let dir = session(&server, "utf8");
let script = dir.path().join("hearts.sh");
std::fs::write(
&script,
format!("#!/bin/sh\nprintf '\\342\\235\\244%.0s' $(seq {HEARTS})\nprintf ENDOFRUN\n"),
)
.unwrap();
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap();
let mut client = ControlClient::attach(&server.tmux, "utf8", DEFAULT_SIZE).unwrap();
let _ = client.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::SessionChanged { .. })
});
let pane = server.tmux.list_panes("utf8").unwrap().remove(0);
client
.send_command(&format!("send-keys -t {pane} -l {}", script.display()))
.unwrap();
client
.send_command(&format!("send-keys -t {pane} Enter"))
.unwrap();
let mut seen: Vec<u8> = Vec::new();
let _ = client.wait_for(TIMEOUT, |event| {
if let control::Event::Output { data, .. } = event {
seen.extend_from_slice(data);
}
seen.windows(8).any(|w| w == b"ENDOFRUN")
});
let text = String::from_utf8_lossy(&seen);
assert!(
!text.contains('\u{fffd}'),
"no byte may be replaced in transit; saw {} replacement characters",
text.matches('\u{fffd}').count()
);
assert!(
text.matches('\u{2764}').count() >= HEARTS,
"wanted {HEARTS} hearts, saw {}",
text.matches('\u{2764}').count()
);
}
#[test]
fn killing_the_session_ends_the_connection() {
let server = TestServer::new();
let _dir = session(&server, "gone");
let client = ControlClient::attach(&server.tmux, "gone", DEFAULT_SIZE).unwrap();
let _ = client.wait_for(TIMEOUT, |e| {
matches!(e, control::Event::SessionChanged { .. })
});
server.tmux.kill_session("gone").unwrap();
let ended = client
.wait_for(TIMEOUT, |e| matches!(e, control::Event::Exit { .. }))
.is_ok()
|| matches!(client.try_event(), Err(TryRecvError::Disconnected));
assert!(ended, "the client should notice its session vanish");
}
#[test]
fn shutdown_is_idempotent() {
let server = TestServer::new();
let _dir = session(&server, "bye");
let mut client = ControlClient::attach(&server.tmux, "bye", DEFAULT_SIZE).unwrap();
client.shutdown().unwrap();
client.shutdown().unwrap();
}
}