| 1 | //! A second tmux client, in control mode, used as a query and event channel. |
| 2 | //! |
| 3 | //! The interactive client in the pty cannot answer questions — it is busy being |
| 4 | //! a terminal. Control mode gives us a client that speaks a line protocol |
| 5 | //! instead of drawing: we write commands, tmux answers in `%begin`/`%end` |
| 6 | //! blocks, and pushes `%`-prefixed notifications whenever the server changes. |
| 7 | //! That replaces both halves of the obvious alternative: no `tmux` process |
| 8 | //! spawned per poll, and no polling at all for things tmux will tell us about. |
| 9 | //! |
| 10 | //! The client attaches with `-f read-only,ignore-size,no-output`: |
| 11 | //! |
| 12 | //! - `read-only` — it can never send keystrokes to a pane. |
| 13 | //! - `ignore-size` — an 80x24 control client would otherwise shrink the |
| 14 | //! session's windows to fit itself, which the user would see immediately. |
| 15 | //! - `no-output` — without it tmux streams every byte every pane produces, to a |
| 16 | //! client that has no use for any of it. |
| 17 | //! |
| 18 | //! `read-only` governs *keys*, not commands, so commands issued here still take |
| 19 | //! effect. Nothing in this module decides to run one: the caller passes a fully |
| 20 | //! built line, and the only lines built from client input are the allowlisted |
| 21 | //! ones in `server.rs`. |
| 22 | |
| 23 | use std::collections::VecDeque; |
| 24 | use std::process::Stdio; |
| 25 | use std::sync::Arc; |
| 26 | |
| 27 | use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; |
| 28 | use tokio::sync::{Mutex, mpsc, oneshot}; |
| 29 | |
| 30 | /// A `%`-prefixed line tmux sent us unprompted, verbatim. |
| 31 | pub type Notification = String; |
| 32 | |
| 33 | type Reply = oneshot::Sender<Result<Vec<String>, String>>; |
| 34 | |
| 35 | pub struct Control { |
| 36 | jobs: mpsc::Sender<(String, Reply)>, |
| 37 | /// Held, not used: the client is spawned with `kill_on_drop`, so the handle |
| 38 | /// living exactly as long as this struct is what ties the extra tmux client |
| 39 | /// to the connection that wanted it. |
| 40 | _child: tokio::process::Child, |
| 41 | } |
| 42 | |
| 43 | impl Control { |
| 44 | /// Attach to `session`. The session must already exist — the caller spawns |
| 45 | /// the interactive client first, which creates it. |
| 46 | /// |
| 47 | /// `global_args` must be the profile's server-selection flags, or this |
| 48 | /// attaches to a different tmux server than the terminal is showing. |
| 49 | pub async fn attach( |
| 50 | session: &str, |
| 51 | global_args: &[String], |
| 52 | ) -> std::io::Result<(Control, mpsc::Receiver<Notification>)> { |
| 53 | let mut child = tokio::process::Command::new("tmux") |
| 54 | .args(global_args) |
| 55 | .args([ |
| 56 | "-C", |
| 57 | "attach", |
| 58 | "-t", |
| 59 | session, |
| 60 | "-f", |
| 61 | "read-only,ignore-size,no-output", |
| 62 | ]) |
| 63 | .stdin(Stdio::piped()) |
| 64 | .stdout(Stdio::piped()) |
| 65 | .stderr(Stdio::null()) |
| 66 | .kill_on_drop(true) |
| 67 | .spawn()?; |
| 68 | |
| 69 | let mut stdin = child.stdin.take().expect("stdin piped"); |
| 70 | let stdout = child.stdout.take().expect("stdout piped"); |
| 71 | |
| 72 | let pending: Arc<Mutex<VecDeque<Reply>>> = Arc::new(Mutex::new(VecDeque::new())); |
| 73 | let (notify_tx, notify_rx) = mpsc::channel::<Notification>(64); |
| 74 | let (jobs_tx, mut jobs_rx) = mpsc::channel::<(String, Reply)>(16); |
| 75 | |
| 76 | // Writer. Queues the reply slot *before* writing, so a fast answer can |
| 77 | // never arrive with the queue still empty. |
| 78 | let queue = Arc::clone(&pending); |
| 79 | tokio::spawn(async move { |
| 80 | while let Some((line, reply)) = jobs_rx.recv().await { |
| 81 | queue.lock().await.push_back(reply); |
| 82 | if stdin |
| 83 | .write_all(format!("{line}\n").as_bytes()) |
| 84 | .await |
| 85 | .is_err() |
| 86 | || stdin.flush().await.is_err() |
| 87 | { |
| 88 | break; |
| 89 | } |
| 90 | } |
| 91 | // Dropping stdin makes tmux print %exit and the child exit, which |
| 92 | // is how the control client goes away when the sidebar does. |
| 93 | }); |
| 94 | |
| 95 | let mut lines = BufReader::new(stdout).lines(); |
| 96 | |
| 97 | // tmux answers the attach itself with an empty block before we have |
| 98 | // asked anything. Swallowing it here, while nothing is in flight, is |
| 99 | // what keeps replies lined up with requests: leave it for the reader |
| 100 | // loop and it consumes the first real request's slot, silently shifting |
| 101 | // every answer by one. The timeout is for a future tmux that stops |
| 102 | // sending it, so we hang for half a second rather than forever. |
| 103 | let _ = tokio::time::timeout(std::time::Duration::from_millis(500), async { |
| 104 | while let Ok(Some(line)) = lines.next_line().await { |
| 105 | if line.starts_with("%end") || line.starts_with("%error") { |
| 106 | break; |
| 107 | } |
| 108 | } |
| 109 | }) |
| 110 | .await; |
| 111 | |
| 112 | // Reader. |
| 113 | let queue = Arc::clone(&pending); |
| 114 | tokio::spawn(async move { |
| 115 | let mut block: Option<Vec<String>> = None; |
| 116 | while let Ok(Some(line)) = lines.next_line().await { |
| 117 | if line.starts_with("%begin") { |
| 118 | block = Some(Vec::new()); |
| 119 | } else if line.starts_with("%end") || line.starts_with("%error") { |
| 120 | let body = block.take().unwrap_or_default(); |
| 121 | let result = if line.starts_with("%error") { |
| 122 | Err(body.join("\n")) |
| 123 | } else { |
| 124 | Ok(body) |
| 125 | }; |
| 126 | // tmux answers the attach itself with an empty block before |
| 127 | // we have asked anything; an unmatched reply is normal and |
| 128 | // is dropped rather than desynchronising the queue. |
| 129 | if let Some(reply) = queue.lock().await.pop_front() { |
| 130 | let _ = reply.send(result); |
| 131 | } |
| 132 | } else if let Some(body) = block.as_mut() { |
| 133 | body.push(line); |
| 134 | } else if line.starts_with('%') && notify_tx.send(line).await.is_err() { |
| 135 | break; |
| 136 | } |
| 137 | } |
| 138 | // tmux is gone: fail every waiter rather than leaving them hanging. |
| 139 | for reply in queue.lock().await.drain(..) { |
| 140 | let _ = reply.send(Err("control client exited".into())); |
| 141 | } |
| 142 | }); |
| 143 | |
| 144 | let control = Control { |
| 145 | jobs: jobs_tx, |
| 146 | _child: child, |
| 147 | }; |
| 148 | Ok((control, notify_rx)) |
| 149 | } |
| 150 | |
| 151 | /// Run one tmux command line and return its output lines. |
| 152 | /// |
| 153 | /// `line` must be a complete, already-quoted tmux command. Anything derived |
| 154 | /// from client input has to be validated by the caller first — this is a |
| 155 | /// command channel to a live tmux server, not a sandbox. |
| 156 | pub async fn run(&self, line: impl Into<String>) -> Result<Vec<String>, String> { |
| 157 | let (tx, rx) = oneshot::channel(); |
| 158 | self.jobs |
| 159 | .send((line.into(), tx)) |
| 160 | .await |
| 161 | .map_err(|_| "control client is gone".to_string())?; |
| 162 | rx.await.unwrap_or_else(|_| Err("no reply".into())) |
| 163 | } |
| 164 | } |
| 165 | |
| 166 | /// Notifications that mean our view of the server is stale. Everything else |
| 167 | /// (`%output`, `%layout-change`, pane modes) changes nothing we display. |
| 168 | pub fn is_interesting(notification: &str) -> bool { |
| 169 | const PREFIXES: [&str; 8] = [ |
| 170 | "%sessions-changed", |
| 171 | "%session-changed", |
| 172 | "%session-renamed", |
| 173 | "%session-window-changed", |
| 174 | "%client-session-changed", |
| 175 | "%window-add", |
| 176 | "%window-close", |
| 177 | "%window-renamed", |
| 178 | ]; |
| 179 | PREFIXES.iter().any(|p| notification.starts_with(p)) |
| 180 | } |
| 181 | |
| 182 | #[cfg(test)] |
| 183 | mod tests { |
| 184 | use super::*; |
| 185 | |
| 186 | #[test] |
| 187 | fn only_server_shape_changes_are_interesting() { |
| 188 | assert!(is_interesting("%client-session-changed /dev/pts/3 $1 work")); |
| 189 | assert!(is_interesting("%window-renamed @4 editor")); |
| 190 | assert!(!is_interesting("%output %2 hello")); |
| 191 | assert!(!is_interesting("%layout-change @1 abcd,80x24,0,0,1")); |
| 192 | } |
| 193 | |
| 194 | /// Needs a tmux on PATH; skipped rather than failed where there isn't one, |
| 195 | /// so the suite still runs in a bare container. |
| 196 | #[tokio::test] |
| 197 | async fn round_trips_a_command_against_real_tmux() { |
| 198 | if !crate::pty::Profile::tmux_available() { |
| 199 | return; |
| 200 | } |
| 201 | let session = "termbridge-control-test"; |
| 202 | let _ = std::process::Command::new("tmux") |
| 203 | .args(["new-session", "-d", "-s", session]) |
| 204 | .status(); |
| 205 | |
| 206 | let Ok((control, _notifications)) = Control::attach(session, &[]).await else { |
| 207 | return; |
| 208 | }; |
| 209 | let out = control |
| 210 | .run("list-sessions -F '#{session_name}'") |
| 211 | .await |
| 212 | .expect("list-sessions succeeds"); |
| 213 | assert!(out.iter().any(|l| l == session)); |
| 214 | |
| 215 | let err = control.run("definitely-not-a-command").await; |
| 216 | assert!(err.is_err(), "tmux errors come back as Err, not output"); |
| 217 | |
| 218 | let _ = std::process::Command::new("tmux") |
| 219 | .args(["kill-session", "-t", session]) |
| 220 | .status(); |
| 221 | } |
| 222 | } |