返回 CodeWhale
pty.rs
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
311 lines RUST