Skip to main content

command_stream/terminal/
capture.rs

1use super::artifacts::{unroll_terminal_frames, write_terminal_artifacts};
2use super::types::{
3    Asciicast, AsciicastEvent, AsciicastHeader, TerminalCapture, TerminalCaptureError,
4    TerminalCaptureOptions, TerminalCursor, TerminalFrame, TerminalInteraction, TerminalResize,
5};
6use portable_pty::{native_pty_system, Child, CommandBuilder, ExitStatus, MasterPty, PtySize};
7use regex::Regex;
8use std::collections::HashMap;
9use std::io::{Read, Write};
10use std::sync::mpsc;
11use std::time::{Duration, Instant};
12
13const ERASE_SCREEN: &[u8] = b"\x1b[2J";
14
15fn elapsed(started: Instant) -> f64 {
16    (started.elapsed().as_secs_f64() * 1_000_000.0).round() / 1_000_000.0
17}
18
19fn trim_trailing_blank(mut lines: Vec<String>) -> Vec<String> {
20    while lines.last().is_some_and(String::is_empty) {
21        lines.pop();
22    }
23    lines
24}
25
26fn frame(parser: &vt100::Parser, started: Instant) -> TerminalFrame {
27    let screen = parser.screen();
28    let (rows, cols) = screen.size();
29    let (cursor_y, cursor_x) = screen.cursor_position();
30    let lines = trim_trailing_blank(screen.rows(0, cols).collect());
31    TerminalFrame {
32        time: elapsed(started),
33        cols,
34        rows,
35        cursor: TerminalCursor {
36            x: cursor_x,
37            y: cursor_y,
38        },
39        alternate: screen.alternate_screen(),
40        screen: lines.clone(),
41        lines,
42    }
43}
44
45fn same_frame(left: &TerminalFrame, right: &TerminalFrame) -> bool {
46    left.cols == right.cols
47        && left.rows == right.rows
48        && left.cursor == right.cursor
49        && left.alternate == right.alternate
50        && left.lines == right.lines
51}
52
53fn append_frame(frames: &mut Vec<TerminalFrame>, parser: &vt100::Parser, started: Instant) {
54    let next = frame(parser, started);
55    if frames
56        .last()
57        .is_none_or(|previous| !same_frame(previous, &next))
58    {
59        frames.push(next);
60    }
61}
62
63fn render_segments(data: &[u8]) -> Vec<&[u8]> {
64    let positions = data
65        .windows(ERASE_SCREEN.len())
66        .enumerate()
67        .filter_map(|(index, window)| (window == ERASE_SCREEN).then_some(index))
68        .collect::<Vec<_>>();
69    if positions.is_empty() {
70        return vec![data];
71    }
72
73    let mut segments = Vec::new();
74    if positions[0] > 0 {
75        segments.push(&data[..positions[0]]);
76    }
77    for (index, position) in positions.iter().enumerate() {
78        let end = positions.get(index + 1).copied().unwrap_or(data.len());
79        segments.push(&data[*position..end]);
80    }
81    segments
82}
83
84fn drain_complete_render_data(pending: &mut Vec<u8>) -> Vec<u8> {
85    let maximum = pending.len().min(ERASE_SCREEN.len() - 1);
86    let pending_length = (1..=maximum)
87        .rev()
88        .find(|length| ERASE_SCREEN.starts_with(&pending[pending.len() - length..]))
89        .unwrap_or(0);
90    pending.drain(..pending.len() - pending_length).collect()
91}
92
93fn record(asciicast: &mut Asciicast, started: Instant, code: &str, data: impl Into<String>) {
94    asciicast.events.push(AsciicastEvent {
95        time: elapsed(started),
96        code: code.into(),
97        data: data.into(),
98    });
99}
100
101fn apply_interaction(
102    interaction: &TerminalInteraction,
103    writer: &mut dyn Write,
104    master: &dyn MasterPty,
105    parser: &mut vt100::Parser,
106    asciicast: &mut Asciicast,
107    started: Instant,
108) -> Result<(), TerminalCaptureError> {
109    if let Some(text) = &interaction.text {
110        writer
111            .write_all(text.as_bytes())
112            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
113        writer
114            .flush()
115            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
116        record(asciicast, started, "i", text.clone());
117    }
118    if let Some(key) = &interaction.key {
119        writer
120            .write_all(key.sequence().as_bytes())
121            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
122        writer
123            .flush()
124            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
125        record(asciicast, started, "i", key.sequence());
126    }
127    if let Some(resize) = interaction.resize {
128        resize_terminal(master, parser, resize)?;
129        record(
130            asciicast,
131            started,
132            "r",
133            format!("{}x{}", resize.cols, resize.rows),
134        );
135    }
136    Ok(())
137}
138
139fn resize_terminal(
140    master: &dyn MasterPty,
141    parser: &mut vt100::Parser,
142    resize: TerminalResize,
143) -> Result<(), TerminalCaptureError> {
144    master
145        .resize(PtySize {
146            rows: resize.rows,
147            cols: resize.cols,
148            pixel_width: 0,
149            pixel_height: 0,
150        })
151        .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
152    parser.set_size(resize.rows, resize.cols);
153    Ok(())
154}
155
156fn asciicast(options: &TerminalCaptureOptions) -> Asciicast {
157    let mut env = HashMap::new();
158    env.insert("SHELL".into(), options.file.clone());
159    env.insert(
160        "TERM".into(),
161        options
162            .env
163            .get("TERM")
164            .cloned()
165            .unwrap_or_else(|| "xterm-256color".into()),
166    );
167    Asciicast {
168        header: AsciicastHeader {
169            version: 2,
170            width: options.cols,
171            height: options.rows,
172            timestamp: chrono::Utc::now().timestamp(),
173            env,
174        },
175        events: Vec::new(),
176    }
177}
178
179fn spawn_reader(mut reader: Box<dyn Read + Send>) -> mpsc::Receiver<Vec<u8>> {
180    let (sender, receiver) = mpsc::channel();
181    std::thread::spawn(move || {
182        let mut buffer = [0_u8; 8192];
183        loop {
184            match reader.read(&mut buffer) {
185                Ok(0) | Err(_) => break,
186                Ok(length) => {
187                    if sender.send(buffer[..length].to_vec()).is_err() {
188                        break;
189                    }
190                }
191            }
192        }
193    });
194    receiver
195}
196
197fn capture_result(
198    status: portable_pty::ExitStatus,
199    output: String,
200    frames: Vec<TerminalFrame>,
201    interaction_count: usize,
202    asciicast: Asciicast,
203) -> TerminalCapture {
204    TerminalCapture {
205        exit_code: status.exit_code() as i32,
206        signal: status.signal().map(str::to_owned),
207        transcript: unroll_terminal_frames(&frames),
208        output,
209        frames,
210        interaction_count,
211        asciicast,
212    }
213}
214
215/// Readiness condition for [`TerminalSession::wait_for`], mirroring the
216/// `after` / `after_regex` vocabulary of [`TerminalInteraction`].
217#[derive(Debug, Clone)]
218pub enum TerminalPattern {
219    Text(String),
220    Regex(Regex),
221}
222
223impl TerminalPattern {
224    pub fn text(value: impl Into<String>) -> Self {
225        Self::Text(value.into())
226    }
227
228    pub fn regex(pattern: &str) -> Result<Self, TerminalCaptureError> {
229        Regex::new(pattern).map(Self::Regex).map_err(|error| {
230            TerminalCaptureError::new(format!("invalid terminal pattern regex: {error}"), None)
231        })
232    }
233
234    fn matches(&self, output: &str) -> bool {
235        match self {
236            Self::Text(value) => output.contains(value),
237            Self::Regex(pattern) => pattern.is_match(output),
238        }
239    }
240}
241
242/// A pseudoterminal that stays open until the caller closes it, so input may be
243/// sent long after the process started.
244pub struct TerminalSession {
245    options: TerminalCaptureOptions,
246    interaction_regexes: Vec<Option<Regex>>,
247    master: Box<dyn MasterPty + Send>,
248    writer: Box<dyn Write + Send>,
249    child: Box<dyn Child + Send + Sync>,
250    receiver: mpsc::Receiver<Vec<u8>>,
251    started: Instant,
252    parser: vt100::Parser,
253    recording: Asciicast,
254    output: String,
255    frames: Vec<TerminalFrame>,
256    pending_render: Vec<u8>,
257    terminal_has_output: bool,
258    interaction_index: usize,
259    last_output: Option<Instant>,
260    dirty: bool,
261    reader_closed: bool,
262    status: Option<ExitStatus>,
263    timed_out: bool,
264    stop_deadline: Option<Instant>,
265}
266
267impl TerminalSession {
268    fn open(options: TerminalCaptureOptions) -> Result<Self, TerminalCaptureError> {
269        if options.file.is_empty() {
270            return Err(TerminalCaptureError::new(
271                "open_terminal requires a file",
272                None,
273            ));
274        }
275        let interaction_regexes = options
276            .interactions
277            .iter()
278            .map(|interaction| {
279                interaction
280                    .after_regex
281                    .as_ref()
282                    .map(|pattern| {
283                        Regex::new(pattern).map_err(|error| {
284                            TerminalCaptureError::new(
285                                format!("invalid terminal interaction regex: {error}"),
286                                None,
287                            )
288                        })
289                    })
290                    .transpose()
291            })
292            .collect::<Result<Vec<_>, _>>()?;
293        let pty = native_pty_system()
294            .openpty(PtySize {
295                rows: options.rows,
296                cols: options.cols,
297                pixel_width: 0,
298                pixel_height: 0,
299            })
300            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
301        let mut command = CommandBuilder::new(&options.file);
302        command.args(&options.args);
303        if let Some(cwd) = &options.cwd {
304            command.cwd(cwd);
305        }
306        command.env(
307            "TERM",
308            options
309                .env
310                .get("TERM")
311                .map_or("xterm-256color", String::as_str),
312        );
313        for (name, value) in &options.env {
314            command.env(name, value);
315        }
316        let child = pty
317            .slave
318            .spawn_command(command)
319            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
320        drop(pty.slave);
321        let reader = pty
322            .master
323            .try_clone_reader()
324            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
325        let writer = pty
326            .master
327            .take_writer()
328            .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
329        let receiver = spawn_reader(reader);
330        let recording = asciicast(&options);
331        let parser = vt100::Parser::new(options.rows, options.cols, 100_000);
332        Ok(Self {
333            interaction_regexes,
334            master: pty.master,
335            writer,
336            child,
337            receiver,
338            started: Instant::now(),
339            parser,
340            recording,
341            output: String::new(),
342            frames: Vec::new(),
343            pending_render: Vec::new(),
344            terminal_has_output: false,
345            interaction_index: 0,
346            last_output: None,
347            dirty: false,
348            reader_closed: false,
349            status: None,
350            timed_out: false,
351            stop_deadline: None,
352            options,
353        })
354    }
355
356    /// Raw PTY output seen so far.
357    pub fn output(&self) -> &str {
358        &self.output
359    }
360
361    /// Settled states retained so far.
362    pub fn frames(&self) -> &[TerminalFrame] {
363        &self.frames
364    }
365
366    /// Unrolled transcript of the states retained so far.
367    pub fn transcript(&self) -> String {
368        unroll_terminal_frames(&self.frames)
369    }
370
371    /// Whether the child is still alive.
372    pub fn running(&self) -> bool {
373        self.status.is_none()
374    }
375
376    fn read_available(&mut self) {
377        match self.receiver.recv_timeout(Duration::from_millis(5)) {
378            Ok(data) => {
379                let text = String::from_utf8_lossy(&data);
380                self.output.push_str(&text);
381                record(&mut self.recording, self.started, "o", text.into_owned());
382                self.pending_render.extend_from_slice(&data);
383                let render_data = drain_complete_render_data(&mut self.pending_render);
384                let segments = render_segments(&render_data);
385                let segment_count = segments.len();
386                if self.terminal_has_output && render_data.starts_with(ERASE_SCREEN) {
387                    append_frame(&mut self.frames, &self.parser, self.started);
388                }
389                for (index, segment) in segments.into_iter().enumerate() {
390                    self.parser.process(segment);
391                    self.terminal_has_output |= !segment.is_empty();
392                    if index + 1 < segment_count {
393                        append_frame(&mut self.frames, &self.parser, self.started);
394                    }
395                }
396                self.last_output = Some(Instant::now());
397                self.dirty = true;
398                if self
399                    .options
400                    .stop_marker
401                    .as_ref()
402                    .is_some_and(|marker| self.output.contains(marker))
403                    && self.stop_deadline.is_none()
404                {
405                    append_frame(&mut self.frames, &self.parser, self.started);
406                    self.stop_deadline = Some(Instant::now() + self.options.stop_marker_grace);
407                }
408            }
409            Err(mpsc::RecvTimeoutError::Disconnected) => self.reader_closed = true,
410            Err(mpsc::RecvTimeoutError::Timeout) => {}
411        }
412    }
413
414    fn idle_for(&self) -> Duration {
415        self.last_output
416            .map_or_else(|| self.started.elapsed(), |instant| instant.elapsed())
417    }
418
419    fn apply_scripted_interactions(&mut self) -> Result<(), TerminalCaptureError> {
420        while let Some(interaction) = self.options.interactions.get(self.interaction_index) {
421            if interaction
422                .after
423                .as_ref()
424                .is_some_and(|marker| !self.output.contains(marker))
425            {
426                break;
427            }
428            if self.interaction_regexes[self.interaction_index]
429                .as_ref()
430                .is_some_and(|pattern| !pattern.is_match(&self.output))
431            {
432                break;
433            }
434            if interaction.idle_duration > Duration::ZERO
435                && self.idle_for() < interaction.idle_duration
436            {
437                break;
438            }
439            let interaction = interaction.clone();
440            append_frame(&mut self.frames, &self.parser, self.started);
441            apply_interaction(
442                &interaction,
443                self.writer.as_mut(),
444                self.master.as_ref(),
445                &mut self.parser,
446                &mut self.recording,
447                self.started,
448            )?;
449            self.interaction_index += 1;
450        }
451        Ok(())
452    }
453
454    /// Advance the capture by one step: read pending output, apply any scripted
455    /// interaction whose readiness condition became true, retain settled states,
456    /// and enforce the optional deadline.
457    fn poll(&mut self) -> Result<(), TerminalCaptureError> {
458        self.read_available();
459        self.apply_scripted_interactions()?;
460
461        if self.dirty
462            && self
463                .last_output
464                .is_some_and(|instant| instant.elapsed() >= self.options.settle_duration)
465        {
466            append_frame(&mut self.frames, &self.parser, self.started);
467            self.dirty = false;
468        }
469        if self.status.is_none() {
470            self.status = self
471                .child
472                .try_wait()
473                .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?;
474        }
475        if self.status.is_none() {
476            let expired = self
477                .options
478                .timeout
479                .is_some_and(|timeout| self.started.elapsed() >= timeout);
480            let stopped = self
481                .stop_deadline
482                .is_some_and(|deadline| Instant::now() >= deadline);
483            if expired || stopped {
484                self.timed_out = expired;
485                self.stop()?;
486            }
487        }
488        Ok(())
489    }
490
491    fn finished(&self) -> bool {
492        self.status.is_some() && self.reader_closed
493    }
494
495    fn stop(&mut self) -> Result<(), TerminalCaptureError> {
496        let _ = self.child.kill();
497        self.status = Some(
498            self.child
499                .wait()
500                .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?,
501        );
502        Ok(())
503    }
504
505    /// Block until `pattern` has been seen and, when `idle` is non-zero, no
506    /// further output arrived for that long. New output restarts the idle wait.
507    pub fn wait_for(
508        &mut self,
509        pattern: &TerminalPattern,
510        idle: Duration,
511        timeout: Option<Duration>,
512    ) -> Result<(), TerminalCaptureError> {
513        let deadline = timeout.map(|limit| Instant::now() + limit);
514        loop {
515            if pattern.matches(&self.output) && self.idle_for() >= idle {
516                return Ok(());
517            }
518            if self.status.is_some() {
519                self.poll()?;
520                if pattern.matches(&self.output) {
521                    return Ok(());
522                }
523                if self.finished() {
524                    return Err(TerminalCaptureError::new(
525                        "terminal exited before the expected output arrived",
526                        None,
527                    ));
528                }
529                continue;
530            }
531            if deadline.is_some_and(|limit| Instant::now() >= limit) {
532                return Err(TerminalCaptureError::new(
533                    format!(
534                        "terminal wait_for timed out after {} ms",
535                        timeout.unwrap_or_default().as_millis()
536                    ),
537                    None,
538                ));
539            }
540            self.poll()?;
541        }
542    }
543
544    /// Send text, a named key, or a resize to the live terminal, using the same
545    /// vocabulary as [`TerminalInteraction`].
546    pub fn send(&mut self, interaction: &TerminalInteraction) -> Result<(), TerminalCaptureError> {
547        if let Some(marker) = &interaction.after {
548            let pattern = TerminalPattern::text(marker.clone());
549            self.wait_for(&pattern, interaction.idle_duration, None)?;
550        } else if let Some(expression) = &interaction.after_regex {
551            let pattern = TerminalPattern::regex(expression)?;
552            self.wait_for(&pattern, interaction.idle_duration, None)?;
553        } else if interaction.idle_duration > Duration::ZERO {
554            while self.idle_for() < interaction.idle_duration && self.status.is_none() {
555                self.poll()?;
556            }
557        }
558        if self.status.is_some() {
559            return Err(TerminalCaptureError::new(
560                "terminal session has already exited",
561                None,
562            ));
563        }
564        append_frame(&mut self.frames, &self.parser, self.started);
565        apply_interaction(
566            interaction,
567            self.writer.as_mut(),
568            self.master.as_ref(),
569            &mut self.parser,
570            &mut self.recording,
571            self.started,
572        )
573    }
574
575    /// Wait for a child that exits on its own, then produce the capture.
576    pub fn finish(mut self) -> Result<TerminalCapture, TerminalCaptureError> {
577        while !self.finished() {
578            self.poll()?;
579        }
580        self.into_capture()
581    }
582
583    /// Stop the child if it is still running, then produce the capture and write
584    /// any configured artifacts.
585    pub fn close(mut self) -> Result<TerminalCapture, TerminalCaptureError> {
586        if self.status.is_none() {
587            self.stop()?;
588        }
589        while !self.finished() {
590            self.read_available();
591        }
592        self.into_capture()
593    }
594
595    fn into_capture(mut self) -> Result<TerminalCapture, TerminalCaptureError> {
596        let pending = std::mem::take(&mut self.pending_render);
597        self.parser.process(&pending);
598        append_frame(&mut self.frames, &self.parser, self.started);
599        let capture = capture_result(
600            self.status
601                .clone()
602                .expect("child status is available after the capture loop"),
603            std::mem::take(&mut self.output),
604            std::mem::take(&mut self.frames),
605            self.interaction_index,
606            std::mem::replace(&mut self.recording, asciicast(&self.options)),
607        );
608        if let Some(directory) = &self.options.artifact_directory {
609            write_terminal_artifacts(
610                directory,
611                &capture.frames,
612                &capture.transcript,
613                &capture.asciicast,
614            )?;
615        }
616        if self.timed_out {
617            return Err(TerminalCaptureError::new(
618                format!(
619                    "terminal command timed out after {} ms",
620                    self.options.timeout.unwrap_or_default().as_millis()
621                ),
622                Some(capture),
623            ));
624        }
625        Ok(capture)
626    }
627}
628
629/// Open a terminal session that stays alive until the caller closes it.
630///
631/// Unlike [`capture_terminal`], input may be sent at any later point through
632/// [`TerminalSession::send`], and readiness can be awaited with
633/// [`TerminalSession::wait_for`]. `options.timeout` defaults to `None` here, so
634/// nothing terminates the child until [`TerminalSession::close`] is called.
635pub fn open_terminal(
636    options: TerminalCaptureOptions,
637) -> Result<TerminalSession, TerminalCaptureError> {
638    TerminalSession::open(TerminalCaptureOptions {
639        timeout: None,
640        ..options
641    })
642}
643
644/// Run a command inside a real pseudoterminal and retain its settled TUI states.
645pub fn capture_terminal(
646    options: TerminalCaptureOptions,
647) -> Result<TerminalCapture, TerminalCaptureError> {
648    if options.file.is_empty() {
649        return Err(TerminalCaptureError::new(
650            "capture_terminal requires a file",
651            None,
652        ));
653    }
654    TerminalSession::open(options)?.finish()
655}
656
657pub async fn capture_terminal_async(
658    options: TerminalCaptureOptions,
659) -> Result<TerminalCapture, TerminalCaptureError> {
660    tokio::task::spawn_blocking(move || capture_terminal(options))
661        .await
662        .map_err(|error| TerminalCaptureError::new(error.to_string(), None))?
663}