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#[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
242pub 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 pub fn output(&self) -> &str {
358 &self.output
359 }
360
361 pub fn frames(&self) -> &[TerminalFrame] {
363 &self.frames
364 }
365
366 pub fn transcript(&self) -> String {
368 unroll_terminal_frames(&self.frames)
369 }
370
371 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 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 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 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 pub fn finish(mut self) -> Result<TerminalCapture, TerminalCaptureError> {
577 while !self.finished() {
578 self.poll()?;
579 }
580 self.into_capture()
581 }
582
583 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
629pub fn open_terminal(
636 options: TerminalCaptureOptions,
637) -> Result<TerminalSession, TerminalCaptureError> {
638 TerminalSession::open(TerminalCaptureOptions {
639 timeout: None,
640 ..options
641 })
642}
643
644pub 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}