| 1 | //! Pseudo-terminal session wrapping `portable-pty`. |
| 2 | //! |
| 3 | //! Spawns a binary in a real PTY, pumps the child's stdout into an in-memory |
| 4 | //! buffer on a background thread, and exposes write/wait/kill primitives |
| 5 | //! the test harness composes. |
| 6 | //! |
| 7 | //! The reader thread is necessary because `portable-pty`'s reader is blocking |
| 8 | //! and the test thread must remain free to send input + poll for screen |
| 9 | //! changes. |
| 10 | |
| 11 | use anyhow::{Context, Result}; |
| 12 | use portable_pty::{Child, CommandBuilder, MasterPty, PtySize, native_pty_system}; |
| 13 | use std::io::{Read, Write}; |
| 14 | use std::path::Path; |
| 15 | use std::sync::{Arc, Mutex}; |
| 16 | use std::thread::{self, JoinHandle}; |
| 17 | use std::time::{Duration, Instant}; |
| 18 | |
| 19 | pub struct PtySession { |
| 20 | /// Held (not read) so the PTY master stays open for the child's lifetime. |
| 21 | master: Box<dyn MasterPty + Send>, |
| 22 | child: Box<dyn Child + Send + Sync>, |
| 23 | writer: Box<dyn Write + Send>, |
| 24 | buffer: Arc<Mutex<Vec<u8>>>, |
| 25 | /// Every byte the child ever wrote, never drained. `buffer` is consumed by |
| 26 | /// the frame parser, which is the wrong shape for assertions about the |
| 27 | /// control stream itself — terminal-mode setup/teardown is only visible as |
| 28 | /// escape sequences, and a mode that was enabled and then disabled leaves |
| 29 | /// no trace on the rendered screen at all. |
| 30 | transcript: Arc<Mutex<Vec<u8>>>, |
| 31 | reader_handle: Option<JoinHandle<()>>, |
| 32 | /// The signal that killed the child, when `wait_until` reaped one. |
| 33 | /// `portable_pty` folds a signal death into exit code 1, which reads as a |
| 34 | /// deliberate `exit(1)`; a SIGPIPE death spent a debugging session that way. |
| 35 | signal: Option<String>, |
| 36 | } |
| 37 | |
| 38 | pub struct PtySessionBuilder<'a> { |
| 39 | program: &'a Path, |
| 40 | args: Vec<String>, |
| 41 | cwd: Option<&'a Path>, |
| 42 | env: Vec<(String, String)>, |
| 43 | rows: u16, |
| 44 | cols: u16, |
| 45 | clear_env: bool, |
| 46 | } |
| 47 | |
| 48 | impl<'a> PtySessionBuilder<'a> { |
| 49 | pub fn new(program: &'a Path) -> Self { |
| 50 | Self { |
| 51 | program, |
| 52 | args: Vec::new(), |
| 53 | cwd: None, |
| 54 | env: Vec::new(), |
| 55 | rows: 40, |
| 56 | cols: 120, |
| 57 | clear_env: false, |
| 58 | } |
| 59 | } |
| 60 | |
| 61 | pub fn args<I, S>(mut self, args: I) -> Self |
| 62 | where |
| 63 | I: IntoIterator<Item = S>, |
| 64 | S: Into<String>, |
| 65 | { |
| 66 | self.args.extend(args.into_iter().map(Into::into)); |
| 67 | self |
| 68 | } |
| 69 | |
| 70 | pub fn cwd(mut self, p: &'a Path) -> Self { |
| 71 | self.cwd = Some(p); |
| 72 | self |
| 73 | } |
| 74 | |
| 75 | pub fn env(mut self, k: impl Into<String>, v: impl Into<String>) -> Self { |
| 76 | self.env.push((k.into(), v.into())); |
| 77 | self |
| 78 | } |
| 79 | |
| 80 | /// Wipe the inherited environment before applying explicit `env(..)` |
| 81 | /// overrides. Use for sealed scenarios that must not see the developer's |
| 82 | /// real `~/.deepseek/`, `$HOME`, or API keys. |
| 83 | pub fn clear_env(mut self, yes: bool) -> Self { |
| 84 | self.clear_env = yes; |
| 85 | self |
| 86 | } |
| 87 | |
| 88 | pub fn size(mut self, rows: u16, cols: u16) -> Self { |
| 89 | self.rows = rows; |
| 90 | self.cols = cols; |
| 91 | self |
| 92 | } |
| 93 | |
| 94 | pub fn spawn(self) -> Result<PtySession> { |
| 95 | let pty_system = native_pty_system(); |
| 96 | let pair = pty_system |
| 97 | .openpty(PtySize { |
| 98 | rows: self.rows, |
| 99 | cols: self.cols, |
| 100 | pixel_width: 0, |
| 101 | pixel_height: 0, |
| 102 | }) |
| 103 | .context("openpty")?; |
| 104 | |
| 105 | let mut cmd = CommandBuilder::new(self.program); |
| 106 | for a in &self.args { |
| 107 | cmd.arg(a); |
| 108 | } |
| 109 | if let Some(cwd) = self.cwd { |
| 110 | cmd.cwd(cwd); |
| 111 | } |
| 112 | if self.clear_env { |
| 113 | cmd.env_clear(); |
| 114 | if let Some(path) = std::env::var_os("PATH") { |
| 115 | cmd.env("PATH", path); |
| 116 | } |
| 117 | } |
| 118 | // TERM must be set to something xterm-ish so crossterm enables the |
| 119 | // capabilities the TUI assumes (256 color, bracketed paste, …). |
| 120 | cmd.env("TERM", "xterm-256color"); |
| 121 | cmd.env("COLORTERM", "truecolor"); |
| 122 | for (k, v) in &self.env { |
| 123 | cmd.env(k, v); |
| 124 | } |
| 125 | |
| 126 | let child = pair.slave.spawn_command(cmd).context("spawn child")?; |
| 127 | // Drop the slave end so EOF propagates correctly when the child exits. |
| 128 | drop(pair.slave); |
| 129 | |
| 130 | let mut reader = pair.master.try_clone_reader().context("clone reader")?; |
| 131 | let writer = pair.master.take_writer().context("take writer")?; |
| 132 | |
| 133 | let buffer: Arc<Mutex<Vec<u8>>> = Arc::new(Mutex::new(Vec::new())); |
| 134 | let transcript: Arc<Mutex<Vec<u8>>> = Arc::new(Mutex::new(Vec::new())); |
| 135 | let buf_thread = Arc::clone(&buffer); |
| 136 | let transcript_thread = Arc::clone(&transcript); |
| 137 | let reader_handle = thread::Builder::new() |
| 138 | .name("qa-pty-reader".into()) |
| 139 | .spawn(move || { |
| 140 | let mut chunk = [0u8; 8192]; |
| 141 | loop { |
| 142 | match reader.read(&mut chunk) { |
| 143 | Ok(0) => break, |
| 144 | Ok(n) => { |
| 145 | if let Ok(mut b) = buf_thread.lock() { |
| 146 | b.extend_from_slice(&chunk[..n]); |
| 147 | } |
| 148 | if let Ok(mut t) = transcript_thread.lock() { |
| 149 | t.extend_from_slice(&chunk[..n]); |
| 150 | } |
| 151 | } |
| 152 | Err(_) => break, |
| 153 | } |
| 154 | } |
| 155 | }) |
| 156 | .context("reader thread")?; |
| 157 | |
| 158 | Ok(PtySession { |
| 159 | master: pair.master, |
| 160 | child, |
| 161 | writer, |
| 162 | buffer, |
| 163 | transcript, |
| 164 | reader_handle: Some(reader_handle), |
| 165 | signal: None, |
| 166 | }) |
| 167 | } |
| 168 | } |
| 169 | |
| 170 | impl PtySession { |
| 171 | pub fn builder(program: &Path) -> PtySessionBuilder<'_> { |
| 172 | PtySessionBuilder::new(program) |
| 173 | } |
| 174 | |
| 175 | pub fn pid(&self) -> Option<u32> { |
| 176 | self.child.process_id() |
| 177 | } |
| 178 | |
| 179 | /// The signal that killed the child, once it has been reaped. |
| 180 | pub fn signal(&self) -> Option<&str> { |
| 181 | self.signal.as_deref() |
| 182 | } |
| 183 | |
| 184 | pub fn write_bytes(&mut self, bytes: &[u8]) -> Result<()> { |
| 185 | self.writer.write_all(bytes).context("pty write")?; |
| 186 | self.writer.flush().context("pty flush")?; |
| 187 | Ok(()) |
| 188 | } |
| 189 | |
| 190 | pub fn resize(&self, rows: u16, cols: u16) -> Result<()> { |
| 191 | self.master |
| 192 | .resize(PtySize { |
| 193 | rows, |
| 194 | cols, |
| 195 | pixel_width: 0, |
| 196 | pixel_height: 0, |
| 197 | }) |
| 198 | .context("pty resize") |
| 199 | } |
| 200 | |
| 201 | /// Drain any bytes the reader thread has pushed into the buffer. Returns |
| 202 | /// the bytes read this call. Non-blocking — returns immediately even if |
| 203 | /// the buffer is empty. |
| 204 | /// Every byte the child has written so far, including bytes already fed |
| 205 | /// to the frame parser. Non-destructive, so it can be sampled repeatedly. |
| 206 | pub fn transcript(&self) -> Vec<u8> { |
| 207 | self.transcript |
| 208 | .lock() |
| 209 | .unwrap_or_else(|e| e.into_inner()) |
| 210 | .clone() |
| 211 | } |
| 212 | |
| 213 | pub fn drain(&mut self) -> Vec<u8> { |
| 214 | let mut b = self.buffer.lock().unwrap_or_else(|e| e.into_inner()); |
| 215 | std::mem::take(&mut *b) |
| 216 | } |
| 217 | |
| 218 | /// Block until the child exits or the deadline passes. Returns the exit |
| 219 | /// status if reaped, or `None` on timeout. |
| 220 | pub fn wait_until(&mut self, deadline: Instant) -> Option<i32> { |
| 221 | loop { |
| 222 | match self.child.try_wait() { |
| 223 | Ok(Some(status)) => { |
| 224 | self.signal = status.signal().map(str::to_owned); |
| 225 | return Some(status.exit_code() as i32); |
| 226 | } |
| 227 | Ok(None) => {} |
| 228 | Err(_) => return None, |
| 229 | } |
| 230 | if Instant::now() >= deadline { |
| 231 | return None; |
| 232 | } |
| 233 | thread::sleep(Duration::from_millis(20)); |
| 234 | } |
| 235 | } |
| 236 | |
| 237 | /// Send SIGTERM-equivalent and wait briefly. Returns the exit status if |
| 238 | /// the child reaped within `grace`, or `None` otherwise. |
| 239 | pub fn shutdown(mut self, grace: Duration) -> Option<i32> { |
| 240 | self.kill_and_join_reader(grace) |
| 241 | } |
| 242 | |
| 243 | fn kill_and_join_reader(&mut self, grace: Duration) -> Option<i32> { |
| 244 | // Name the teardown for the watchdog: a wedge here used to be the whole |
| 245 | // bug, so "teardown: kill child" is the message worth seeing. |
| 246 | super::watchdog::progress("teardown: kill child + reap group"); |
| 247 | let _ = self.child.kill(); |
| 248 | // Killing only the direct child is not enough: the TUI spawns shells, |
| 249 | // and a descendant that escaped into its own session keeps the PTY |
| 250 | // slave open, so the reader never sees EOF. Reap the whole group. |
| 251 | self.kill_process_group(); |
| 252 | let exit = self.wait_until(Instant::now() + grace); |
| 253 | if let Some(handle) = self.reader_handle.take() { |
| 254 | join_reader_bounded(handle); |
| 255 | } |
| 256 | exit |
| 257 | } |
| 258 | |
| 259 | /// SIGKILL the child's process group, best effort. |
| 260 | /// |
| 261 | /// `portable_pty`'s `Child::kill` signals one pid. A grandchild in its own |
| 262 | /// session survives it and holds the inherited slave fd, which is the state |
| 263 | /// that made the reader join below unbounded. |
| 264 | #[cfg(unix)] |
| 265 | fn kill_process_group(&mut self) { |
| 266 | let Some(pid) = self.child.process_id() else { |
| 267 | return; |
| 268 | }; |
| 269 | let Ok(pid) = i32::try_from(pid) else { |
| 270 | return; |
| 271 | }; |
| 272 | // SAFETY: `killpg` on a pid we spawned; a stale pid returns ESRCH |
| 273 | // rather than signalling an unrelated group, because the child has not |
| 274 | // been reaped yet at this point. |
| 275 | unsafe { |
| 276 | libc::killpg(pid, libc::SIGKILL); |
| 277 | } |
| 278 | } |
| 279 | |
| 280 | #[cfg(not(unix))] |
| 281 | fn kill_process_group(&mut self) {} |
| 282 | } |
| 283 | |
| 284 | /// Bounded join for the PTY reader thread. |
| 285 | /// |
| 286 | /// The previous code said "don't block forever" but called `handle.join()`, |
| 287 | /// which does exactly that when a descendant still holds the PTY slave open — |
| 288 | /// `read()` never returns EOF. Because libtest has no per-test timeout, that |
| 289 | /// turned any *failing* PTY test into an infinite hang: the assertion returns |
| 290 | /// `Err`, `?` drops the harness, and the drop blocks forever. Hand the join to |
| 291 | /// a helper thread and move on; the reader exits on its own once the pipe |
| 292 | /// finally closes, and the process exits at the end of the test binary anyway. |
| 293 | /// Same shape as `READER_JOIN_GRACE` in `tools/shell.rs` (#52). |
| 294 | fn join_reader_bounded(handle: JoinHandle<()>) { |
| 295 | const READER_JOIN_GRACE: Duration = Duration::from_secs(2); |
| 296 | let (done_tx, done_rx) = std::sync::mpsc::channel(); |
| 297 | let _ = thread::Builder::new() |
| 298 | .name("qa-pty-reader-join".into()) |
| 299 | .spawn(move || { |
| 300 | let _ = handle.join(); |
| 301 | let _ = done_tx.send(()); |
| 302 | }); |
| 303 | let _ = done_rx.recv_timeout(READER_JOIN_GRACE); |
| 304 | } |
| 305 | |
| 306 | impl Drop for PtySession { |
| 307 | fn drop(&mut self) { |
| 308 | let _ = self.kill_and_join_reader(Duration::from_secs(2)); |
| 309 | } |
| 310 | } |
| 311 |