1use std::collections::VecDeque;
20use std::ffi::OsString;
21use std::io::{self, Read};
22use std::path::PathBuf;
23use std::process::{Child, Command, Stdio};
24use std::sync::mpsc::{self, RecvTimeoutError, SyncSender};
25use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
26use std::time::{Duration, Instant};
27
28const POLL: Duration = Duration::from_millis(10);
30
31pub(super) const CHUNK: usize = 4096;
33
34const MAX_LINE: usize = 64 * 1024;
37
38const QUEUE: usize = 1024;
42
43#[derive(Debug, Clone, PartialEq, Eq)]
59pub struct Process {
60 program: OsString,
61 args: Vec<OsString>,
62 dir: Option<PathBuf>,
63 env: Vec<(OsString, OsString)>,
64 pty: Option<(u16, u16)>,
65 no_stdin: bool,
66 cleared: bool,
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
73pub enum Line {
74 Out(String),
76 Err(String),
78}
79
80#[derive(Debug, Clone, PartialEq, Eq)]
82pub enum ProcessOutcome {
83 Finished {
85 code: Option<i32>,
87 },
88 Cancelled,
90}
91
92#[derive(Debug, Clone, Copy, PartialEq, Eq)]
100#[non_exhaustive]
101pub struct Keep {
102 bytes: usize,
103 lines: Option<usize>,
104 limit: Option<Duration>,
105}
106
107impl Keep {
108 #[must_use]
113 pub fn bytes(bytes: usize) -> Self {
114 Self { bytes, lines: None, limit: None }
115 }
116
117 #[must_use]
120 pub fn lines(mut self, lines: usize) -> Self {
121 self.lines = Some(lines);
122 self
123 }
124
125 #[must_use]
129 pub fn limit(mut self, limit: Duration) -> Self {
130 self.limit = Some(limit);
131 self
132 }
133}
134
135#[derive(Debug, Clone, PartialEq, Eq)]
137#[non_exhaustive]
138pub struct Collected {
139 pub text: String,
143 pub outcome: ProcessOutcome,
146 pub trimmed: bool,
148 pub timed_out: bool,
150 pub cancelled: bool,
152}
153
154impl Process {
155 #[must_use]
157 pub fn new(program: impl Into<OsString>) -> Self {
158 Self {
159 program: program.into(),
160 args: Vec::new(),
161 dir: None,
162 env: Vec::new(),
163 pty: None,
164 no_stdin: false,
165 cleared: false,
166 }
167 }
168
169 #[must_use]
171 pub fn arg(mut self, arg: impl Into<OsString>) -> Self {
172 self.args.push(arg.into());
173 self
174 }
175
176 #[must_use]
178 pub fn args(mut self, args: impl IntoIterator<Item = impl Into<OsString>>) -> Self {
179 self.args.extend(args.into_iter().map(Into::into));
180 self
181 }
182
183 #[must_use]
185 pub fn dir(mut self, dir: impl Into<PathBuf>) -> Self {
186 self.dir = Some(dir.into());
187 self
188 }
189
190 #[must_use]
193 pub fn env(mut self, key: impl Into<OsString>, value: impl Into<OsString>) -> Self {
194 self.env.push((key.into(), value.into()));
195 self
196 }
197
198 #[must_use]
222 pub fn clear_env(mut self) -> Self {
223 self.cleared = true;
224 self
225 }
226
227 #[must_use]
236 pub fn pty(mut self, cols: u16, rows: u16) -> Self {
237 self.pty = Some((cols, rows));
238 self
239 }
240
241 #[must_use]
253 pub fn no_stdin(mut self) -> Self {
254 self.no_stdin = true;
255 self
256 }
257
258 pub fn run(self, cancel: &dyn Fn() -> bool, on_line: &mut dyn FnMut(Line)) -> io::Result<ProcessOutcome> {
291 self.run_inner(cancel, on_line, None)
292 }
293
294 pub fn run_with_overwritten(
325 self,
326 cancel: &dyn Fn() -> bool,
327 on_line: &mut dyn FnMut(Line),
328 on_overwritten: &mut dyn FnMut(Line),
329 ) -> io::Result<ProcessOutcome> {
330 self.run_inner(cancel, on_line, Some(on_overwritten))
331 }
332
333 fn run_inner(
336 self,
337 cancel: &dyn Fn() -> bool,
338 on_line: &mut dyn FnMut(Line),
339 mut on_overwritten: Option<&mut dyn FnMut(Line)>,
340 ) -> io::Result<ProcessOutcome> {
341 let frames = on_overwritten.is_some();
342 let group = self.no_stdin && cfg!(unix);
344 let command = self.command(group);
345 let (sender, receiver) = mpsc::sync_channel(QUEUE);
346 let mut child = match self.pty {
347 Some(size) => spawn_on_pty(command, size, &sender, frames, group)?,
348 None => spawn_on_pipes(command, &sender, frames, group)?,
349 };
350 drop(sender);
352 loop {
353 if cancel() {
354 kill(&mut child, group);
355 return Ok(ProcessOutcome::Cancelled);
356 }
357 match receiver.recv_timeout(POLL) {
358 Ok(Sent::Line(line)) => on_line(line),
359 Ok(Sent::Overwritten(frame)) => {
360 if let Some(on_overwritten) = on_overwritten.as_deref_mut() {
361 on_overwritten(frame);
362 }
363 }
364 Err(RecvTimeoutError::Timeout) => {}
365 Err(RecvTimeoutError::Disconnected) => break,
366 }
367 }
368 loop {
371 if let Some(status) = child.try_wait()? {
372 return Ok(ProcessOutcome::Finished { code: status.code() });
373 }
374 if cancel() {
375 kill(&mut child, group);
376 return Ok(ProcessOutcome::Cancelled);
377 }
378 std::thread::sleep(POLL);
379 }
380 }
381
382 pub fn collect(self, keep: Keep, cancel: &dyn Fn() -> bool) -> io::Result<Collected> {
415 let Keep { bytes, lines, limit } = keep;
416 let group = self.no_stdin && cfg!(unix);
417 let command = self.command(group);
418 let merged = Arc::new(Mutex::new(Merged { tail: Tail::new(bytes, lines), open: true }));
419 let mut child = match self.pty {
420 Some(size) => spawn_merged_on_pty(command, size, group, &merged)?,
421 None => spawn_merged(command, group, &merged)?,
422 };
423 let deadline = limit.map(|limit| Instant::now() + limit);
424 let (outcome, timed_out, cancelled) = loop {
428 if cancel() {
429 kill(&mut child, group);
430 break (ProcessOutcome::Cancelled, false, true);
431 }
432 if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
433 kill(&mut child, group);
434 break (ProcessOutcome::Finished { code: None }, true, false);
435 }
436 if !lock(&merged).open
437 && let Some(status) = child.try_wait()?
438 {
439 break (ProcessOutcome::Finished { code: status.code() }, false, false);
440 }
441 std::thread::sleep(POLL);
442 };
443 let mut merged = lock(&merged);
444 Ok(Collected { text: merged.tail.text(), trimmed: merged.tail.trimmed, outcome, timed_out, cancelled })
445 }
446
447 fn command(&self, group: bool) -> Command {
450 let mut command = Command::new(&self.program);
451 command.args(&self.args);
452 if self.no_stdin {
453 command.stdin(Stdio::null());
454 } else {
455 command.stdin(Stdio::inherit());
456 }
457 #[cfg(unix)]
458 if group {
459 use std::os::unix::process::CommandExt;
460 command.process_group(0);
461 }
462 #[cfg(not(unix))]
463 let _ = group;
464 if let Some(dir) = &self.dir {
465 command.current_dir(dir);
466 }
467 if self.cleared {
468 command.env_clear();
469 }
470 for (key, value) in &self.env {
471 command.env(key, value);
472 }
473 command
474 }
475}
476
477enum Sent {
480 Line(Line),
481 Overwritten(Line),
482}
483
484fn kill(child: &mut Child, group: bool) {
487 #[cfg(unix)]
488 if group {
489 let leader = rustix::process::Pid::from_child(child);
490 let _ = rustix::process::kill_process_group(leader, rustix::process::Signal::KILL);
493 }
494 #[cfg(not(unix))]
495 let _ = group;
496 let _ = child.kill();
500 let _ = child.wait();
501}
502
503fn spawn_on_pipes(mut command: Command, sender: &SyncSender<Sent>, frames: bool, group: bool) -> io::Result<Child> {
506 command.stdout(Stdio::piped()).stderr(Stdio::piped());
507 let mut child = command.spawn()?;
508 drop(command);
509 let taken = child.stdout.take().zip(child.stderr.take());
510 let started = match taken {
511 Some((out, err)) => spawn_reader("out", out, Line::Out, frames, sender.clone())
512 .and_then(|()| spawn_reader("err", err, Line::Err, frames, sender.clone())),
513 None => Err(io::Error::other("the child was started without its pipes")),
514 };
515 match started {
516 Ok(()) => Ok(child),
517 Err(error) => {
518 kill(&mut child, group);
519 Err(error)
520 }
521 }
522}
523
524#[cfg(unix)]
527fn spawn_on_pty(
528 mut command: Command,
529 size: (u16, u16),
530 sender: &SyncSender<Sent>,
531 frames: bool,
532 group: bool,
533) -> io::Result<Child> {
534 let terminal = open_pty(&mut command, size)?;
535 let mut child = command.spawn()?;
536 drop(command);
539 match spawn_reader("pty", terminal, Line::Out, frames, sender.clone()) {
540 Ok(()) => Ok(child),
541 Err(error) => {
542 kill(&mut child, group);
543 Err(error)
544 }
545 }
546}
547
548#[cfg(not(unix))]
551fn spawn_on_pty(
552 mut command: Command,
553 size: (u16, u16),
554 _sender: &SyncSender<Sent>,
555 _frames: bool,
556 _group: bool,
557) -> io::Result<Child> {
558 open_pty(&mut command, size)
559}
560
561fn spawn_merged(mut command: Command, group: bool, merged: &Shared) -> io::Result<Child> {
564 let (reader, writer) = io::pipe()?;
565 command.stdout(writer.try_clone()?).stderr(writer);
566 let mut child = command.spawn()?;
567 drop(command);
570 match spawn_collector(reader, Arc::clone(merged)) {
571 Ok(()) => Ok(child),
572 Err(error) => {
573 kill(&mut child, group);
574 Err(error)
575 }
576 }
577}
578
579#[cfg(unix)]
582fn spawn_merged_on_pty(mut command: Command, size: (u16, u16), group: bool, merged: &Shared) -> io::Result<Child> {
583 let terminal = open_pty(&mut command, size)?;
584 let mut child = command.spawn()?;
585 drop(command);
588 match spawn_collector(terminal, Arc::clone(merged)) {
589 Ok(()) => Ok(child),
590 Err(error) => {
591 kill(&mut child, group);
592 Err(error)
593 }
594 }
595}
596
597#[cfg(not(unix))]
599fn spawn_merged_on_pty(mut command: Command, size: (u16, u16), _group: bool, _merged: &Shared) -> io::Result<Child> {
600 open_pty(&mut command, size)
601}
602
603#[cfg(unix)]
606fn open_pty(command: &mut Command, (cols, rows): (u16, u16)) -> io::Result<std::fs::File> {
607 use std::fs::File;
608 use std::os::fd::OwnedFd;
609
610 use rustix::fs::{Mode, OFlags};
611 use rustix::io::{FdFlags, fcntl_setfd};
612 use rustix::pty::{OpenptFlags, grantpt, openpt, ptsname, unlockpt};
613 use rustix::termios::{Winsize, tcsetwinsize};
614
615 #[cfg(any(target_os = "linux", target_os = "android", target_os = "freebsd", target_os = "netbsd"))]
619 let flags = OpenptFlags::RDWR | OpenptFlags::NOCTTY | OpenptFlags::CLOEXEC;
620 #[cfg(not(any(target_os = "linux", target_os = "android", target_os = "freebsd", target_os = "netbsd")))]
621 let flags = OpenptFlags::RDWR | OpenptFlags::NOCTTY;
622 let controller = openpt(flags)?;
623 fcntl_setfd(&controller, FdFlags::CLOEXEC)?;
624 grantpt(&controller)?;
625 unlockpt(&controller)?;
626 tcsetwinsize(&controller, Winsize { ws_row: rows, ws_col: cols, ws_xpixel: 0, ws_ypixel: 0 })?;
627 let name = ptsname(&controller, Vec::new())?;
628 let device: OwnedFd = rustix::fs::open(name, OFlags::RDWR | OFlags::NOCTTY | OFlags::CLOEXEC, Mode::empty())?;
635 command.stdout(Stdio::from(device.try_clone()?)).stderr(Stdio::from(device));
636 Ok(File::from(controller))
637}
638
639#[cfg(not(unix))]
642fn open_pty(_command: &mut Command, _size: (u16, u16)) -> io::Result<std::fs::File> {
643 Err(io::Error::new(io::ErrorKind::Unsupported, "a pseudo-terminal needs a Unix system"))
644}
645
646fn spawn_reader(
649 name: &str,
650 source: impl Read + Send + 'static,
651 tag: fn(String) -> Line,
652 frames: bool,
653 sender: SyncSender<Sent>,
654) -> io::Result<()> {
655 std::thread::Builder::new()
656 .name(format!("quvyta-process-{name}"))
657 .spawn(move || read_lines(source, tag, frames, &sender))
658 .map(|_| ())
659}
660
661fn read_lines(mut source: impl Read, tag: fn(String) -> Line, frames: bool, sender: &SyncSender<Sent>) {
663 let mut chunk = [0_u8; CHUNK];
664 let mut lines = Lines::default();
665 let listening = std::cell::Cell::new(true);
667 let mut on_line = |line| listening.set(listening.get() && sender.send(Sent::Line(tag(line))).is_ok());
668 let mut on_frame = |frame| listening.set(listening.get() && sender.send(Sent::Overwritten(tag(frame))).is_ok());
669 loop {
670 match source.read(&mut chunk) {
671 Ok(0) => break,
672 Ok(count) => {
673 lines.feed_keeping(&chunk[..count], &mut on_line, frames.then_some(&mut on_frame));
674 if !listening.get() {
675 return;
676 }
677 }
678 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
679 Err(_) => break,
682 }
683 }
684 lines.finish_keeping(&mut on_line, frames.then_some(&mut on_frame));
685}
686
687fn spawn_collector(source: impl Read + Send + 'static, merged: Shared) -> io::Result<()> {
690 std::thread::Builder::new()
691 .name("quvyta-process-merged".to_owned())
692 .spawn(move || read_tail(source, &merged))
693 .map(|_| ())
694}
695
696#[derive(Debug)]
700struct Merged {
701 tail: Tail,
702 open: bool,
703}
704
705type Shared = Arc<Mutex<Merged>>;
707
708fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
709 mutex.lock().unwrap_or_else(PoisonError::into_inner)
710}
711
712fn read_tail(mut source: impl Read, merged: &Shared) {
715 let mut chunk = [0_u8; CHUNK];
716 loop {
717 match source.read(&mut chunk) {
718 Ok(0) => break,
719 Ok(count) => lock(merged).tail.feed(&chunk[..count]),
720 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
721 Err(_) => break,
724 }
725 }
726 lock(merged).open = false;
727}
728
729#[derive(Debug)]
733struct Tail {
734 bytes: VecDeque<u8>,
735 lines: usize,
737 ends_line: bool,
739 kept: usize,
741 lines_kept: Option<usize>,
743 trimmed: bool,
745}
746
747impl Tail {
748 fn new(bytes: usize, lines: Option<usize>) -> Self {
750 Self { bytes: VecDeque::new(), lines: 0, ends_line: true, kept: bytes, lines_kept: lines, trimmed: false }
751 }
752
753 fn feed(&mut self, bytes: &[u8]) {
756 let Some(&last) = bytes.last() else {
757 return;
758 };
759 self.ends_line = last == b'\n';
760 self.lines += bytes.iter().filter(|&&byte| byte == b'\n').count();
761 self.bytes.extend(bytes);
762 if let Some(kept) = self.lines_kept {
763 for _ in 0..self.counting().saturating_sub(kept) {
764 self.drop_oldest_line();
765 }
766 }
767 for _ in 0..self.bytes.len().saturating_sub(self.kept) {
768 if self.bytes.pop_front() == Some(b'\n') {
769 self.lines -= 1;
770 }
771 self.trimmed = true;
772 }
773 }
774
775 fn counting(&self) -> usize {
778 self.lines + usize::from(!self.ends_line && !self.bytes.is_empty())
779 }
780
781 fn drop_oldest_line(&mut self) {
783 self.trimmed = true;
784 while let Some(byte) = self.bytes.pop_front() {
785 if byte == b'\n' {
786 self.lines -= 1;
787 break;
788 }
789 }
790 }
791
792 fn text(&mut self) -> String {
795 let bytes = self.bytes.make_contiguous();
796 let start = bytes.iter().copied().take(3).take_while(|byte| byte & 0b1100_0000 == 0b1000_0000).count();
797 String::from_utf8_lossy(&bytes[start..]).into_owned()
798 }
799}
800
801#[derive(Debug, Default)]
804pub(super) struct Lines {
805 buffer: Vec<u8>,
806 pending_return: bool,
808}
809
810impl Lines {
811 pub(super) fn feed(&mut self, bytes: &[u8], emit: &mut impl FnMut(String)) {
813 self.feed_keeping(bytes, emit, None);
814 }
815
816 pub(super) fn feed_keeping(
819 &mut self,
820 bytes: &[u8],
821 emit: &mut impl FnMut(String),
822 mut overwritten: Option<&mut dyn FnMut(String)>,
823 ) {
824 for &byte in bytes {
825 if self.pending_return {
826 match byte {
830 b'\r' => continue,
831 b'\n' => {
832 self.pending_return = false;
833 emit(self.take());
834 continue;
835 }
836 _ => {
837 self.pending_return = false;
838 self.overwrite(&mut overwritten);
839 }
840 }
841 }
842 match byte {
843 b'\r' => self.pending_return = true,
844 b'\n' => emit(self.take()),
845 _ => {
846 self.buffer.push(byte);
847 if self.buffer.len() >= MAX_LINE {
848 self.emit_piece(emit);
849 }
850 }
851 }
852 }
853 }
854
855 fn emit_piece(&mut self, emit: &mut impl FnMut(String)) {
858 let len = self.buffer.len();
861 let mut cut = len;
862 for back in 1..=len.min(3) {
863 let byte = self.buffer[len - back];
864 if byte & 0b1100_0000 != 0b1000_0000 {
865 let width = match byte {
866 0xc0..=0xdf => 2,
867 0xe0..=0xef => 3,
868 0xf0..=0xf7 => 4,
869 _ => 1,
870 };
871 if width > back {
872 cut = len - back;
873 }
874 break;
875 }
876 }
877 let rest = self.buffer.split_off(cut);
878 emit(self.take());
879 self.buffer = rest;
880 }
881
882 fn overwrite(&mut self, overwritten: &mut Option<&mut dyn FnMut(String)>) {
884 match overwritten {
885 Some(overwritten) if !self.buffer.is_empty() => overwritten(self.take()),
886 _ => self.buffer.clear(),
887 }
888 }
889
890 pub(super) fn finish(&mut self, emit: &mut impl FnMut(String)) {
892 self.finish_keeping(emit, None);
893 }
894
895 pub(super) fn finish_keeping(
898 &mut self,
899 emit: &mut impl FnMut(String),
900 mut overwritten: Option<&mut dyn FnMut(String)>,
901 ) {
902 if self.pending_return {
903 self.overwrite(&mut overwritten);
905 self.pending_return = false;
906 }
907 if !self.buffer.is_empty() {
908 emit(self.take());
909 }
910 }
911
912 fn take(&mut self) -> String {
914 let line = String::from_utf8_lossy(&self.buffer).into_owned();
915 self.buffer.clear();
916 line
917 }
918}
919
920#[cfg(test)]
921mod tests {
922 use std::sync::atomic::{AtomicUsize, Ordering};
923 use std::time::Duration;
924
925 use super::{Collected, Keep, Line, Lines, MAX_LINE, Process, ProcessOutcome, Tail};
926
927 fn shell(script: &str) -> (Vec<Line>, ProcessOutcome) {
929 run(Process::new("sh").args(["-c", script]))
930 }
931
932 fn run(process: Process) -> (Vec<Line>, ProcessOutcome) {
934 let mut lines = Vec::new();
935 let outcome = process.run(&|| false, &mut |line| lines.push(line)).expect("the shell starts");
936 (lines, outcome)
937 }
938
939 fn collect(process: Process, keep: Keep) -> Collected {
941 process.collect(keep, &|| false).expect("the shell starts")
942 }
943
944 #[test]
945 fn keeps_the_two_streams_apart_and_reports_the_exit_code() {
946 let (lines, outcome) = shell("echo bir; echo iki >&2; exit 3");
947 assert_eq!(lines.len(), 2, "{lines:?}");
948 assert!(lines.contains(&Line::Out("bir".to_owned())), "{lines:?}");
949 assert!(lines.contains(&Line::Err("iki".to_owned())), "{lines:?}");
950 assert_eq!(outcome, ProcessOutcome::Finished { code: Some(3) });
951 }
952
953 #[test]
954 fn delivers_the_last_line_without_a_newline() {
955 let (lines, outcome) = shell("printf 'son satir'");
956 assert_eq!(lines, vec![Line::Out("son satir".to_owned())]);
957 assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
958 }
959
960 #[test]
961 fn carriage_returns_collapse_into_one_line() {
962 let (lines, _) = shell(r"printf 'a\rbb\rccc\n'");
963 assert_eq!(lines, vec![Line::Out("ccc".to_owned())]);
964 }
965
966 #[test]
967 fn invalid_utf8_becomes_the_replacement_character() {
968 let (lines, _) = shell(r"printf 'a\377b\n'");
969 assert_eq!(lines, vec![Line::Out("a\u{fffd}b".to_owned())]);
970 }
971
972 #[test]
973 fn the_environment_is_inherited_and_one_variable_can_be_replaced() {
974 let (lines, _) = shell("echo ${PATH:+inherited}");
975 assert_eq!(lines, vec![Line::Out("inherited".to_owned())]);
976 let (lines, _) = run(Process::new("sh").args(["-c", "echo $LC_ALL"]).env("LC_ALL", "C"));
977 assert_eq!(lines, vec![Line::Out("C".to_owned())]);
978 }
979
980 #[test]
981 fn a_cleared_environment_gives_the_child_only_what_it_was_told() {
982 let script = r#"printf '%s' "${HOME-unset}""#;
985 let (lines, _) = shell(script);
986 assert_ne!(lines, vec![Line::Out("unset".to_owned())], "the test process really has a home");
987 let path = std::env::var("PATH").expect("the test process was started with a path");
988 let (lines, _) = run(Process::new("sh").args(["-c", script]).clear_env().env("PATH", path.clone()));
989 assert_eq!(lines, vec![Line::Out("unset".to_owned())], "nothing is inherited");
990 let (lines, _) = run(Process::new("sh")
993 .args(["-c", r#"printf '%s' "${PATH-unset}""#])
994 .clear_env()
995 .env("PATH", path.clone()));
996 assert_eq!(lines, vec![Line::Out(path)]);
997 let (lines, _) =
998 run(Process::new("sh").args(["-c", r#"printf '%s' "${EV-unset}""#]).clear_env().env("EV", "1"));
999 assert_eq!(lines, vec![Line::Out("1".to_owned())], "a variable that was given arrives");
1000 }
1001
1002 #[test]
1003 fn runs_in_the_directory_it_is_given() {
1004 let (lines, _) = run(Process::new("sh").args(["-c", "pwd"]).dir("/"));
1005 assert_eq!(lines, vec![Line::Out("/".to_owned())]);
1006 }
1007
1008 #[test]
1009 fn cancelling_kills_a_long_running_child() {
1010 let seen = AtomicUsize::new(0);
1011 let outcome = Process::new("sh")
1012 .args(["-c", "while true; do echo tik; sleep 0.05; done"])
1013 .run(&|| seen.load(Ordering::Relaxed) > 0, &mut |line| {
1014 assert_eq!(line, Line::Out("tik".to_owned()));
1015 seen.fetch_add(1, Ordering::Relaxed);
1016 })
1017 .expect("the shell starts");
1018 assert_eq!(outcome, ProcessOutcome::Cancelled);
1019 assert!(seen.load(Ordering::Relaxed) > 0);
1020 }
1021
1022 #[test]
1023 fn a_missing_program_is_an_error_and_not_a_panic() {
1024 let error = Process::new("quvyta-no-such-program")
1025 .run(&|| false, &mut |_| unreachable!("a missing program writes nothing"))
1026 .expect_err("a missing program cannot run");
1027 assert_eq!(error.kind(), std::io::ErrorKind::NotFound);
1028 }
1029
1030 #[cfg(unix)]
1031 #[test]
1032 fn on_a_pseudo_terminal_the_child_sees_a_terminal_of_the_size_we_gave() {
1033 let (lines, outcome) = run(Process::new("sh").args(["-c", "test -t 1 && stty size <&1"]).pty(100, 24));
1036 assert_eq!(lines, vec![Line::Out("24 100".to_owned())]);
1037 assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1038 }
1039
1040 #[cfg(unix)]
1041 #[test]
1042 fn on_a_pseudo_terminal_both_streams_arrive_as_output() {
1043 let (lines, outcome) = run(Process::new("sh").args(["-c", "echo bir; echo iki >&2"]).pty(80, 24));
1044 assert_eq!(lines, vec![Line::Out("bir".to_owned()), Line::Out("iki".to_owned())]);
1045 assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1046 }
1047
1048 #[cfg(target_os = "linux")]
1050 fn stat_ids(stat: &str) -> [String; 3] {
1051 let fields: Vec<&str> = stat[stat.rfind(')').expect("name") + 2..].split(' ').collect();
1054 [fields[2], fields[3], fields[4]].map(str::to_owned)
1055 }
1056
1057 #[cfg(target_os = "linux")]
1058 #[test]
1059 fn without_stdin_the_child_reads_an_empty_stream() {
1060 let script = r#"readlink /proc/$$/fd/0; read answer; echo "read $?""#;
1061 for process in [Process::new("sh").args(["-c", script]), Process::new("sh").args(["-c", script]).pty(80, 24)] {
1062 let (lines, outcome) = run(process.no_stdin());
1063 assert_eq!(lines, vec![Line::Out("/dev/null".to_owned()), Line::Out("read 1".to_owned())]);
1064 assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1065 }
1066 }
1067
1068 #[cfg(target_os = "linux")]
1069 #[test]
1070 fn only_a_child_without_stdin_gets_a_group_of_its_own_and_it_keeps_the_session() {
1071 let script = "cat /proc/$$/stat";
1072 let ours = stat_ids(&std::fs::read_to_string("/proc/self/stat").expect("stat"));
1073 let ids = |process: Process| {
1074 let (lines, _) = run(process);
1075 let [Line::Out(stat)] = &lines[..] else { panic!("one line: {lines:?}") };
1076 stat_ids(stat)
1077 };
1078 let shared = ids(Process::new("sh").args(["-c", script]));
1079 assert_eq!(shared, ours, "a child reading the terminal stays in the application's group");
1080 for process in [Process::new("sh").args(["-c", script]), Process::new("sh").args(["-c", script]).pty(80, 24)] {
1081 let [group, session, terminal] = ids(process.no_stdin());
1082 assert_ne!(group, ours[0], "a group of its own");
1083 assert_eq!(session, ours[1], "the application's session");
1085 assert_eq!(terminal, ours[2], "the application's controlling terminal");
1086 }
1087 }
1088
1089 #[cfg(target_os = "linux")]
1091 fn ended(pid: &str) -> bool {
1092 std::fs::read_to_string(format!("/proc/{pid}/stat"))
1093 .map_or(true, |stat| stat[stat.rfind(')').expect("name") + 2..].starts_with('Z'))
1094 }
1095
1096 #[cfg(target_os = "linux")]
1097 #[test]
1098 fn cancelling_a_child_without_stdin_ends_the_programs_it_started() {
1099 for pty in [false, true] {
1100 let seen = std::cell::RefCell::new(Vec::new());
1101 let process = Process::new("sh").args(["-c", "sleep 60 & echo $!; sleep 60 & echo $!; wait"]).no_stdin();
1102 let process = if pty { process.pty(80, 24) } else { process };
1103 let outcome = process
1104 .run(&|| seen.borrow().len() == 2, &mut |line| match line {
1105 Line::Out(pid) => seen.borrow_mut().push(pid),
1106 Line::Err(text) => panic!("nothing on standard error: {text}"),
1107 })
1108 .expect("the shell starts");
1109 assert_eq!(outcome, ProcessOutcome::Cancelled);
1110 let pids = seen.into_inner();
1111 let started = std::time::Instant::now();
1112 while !pids.iter().all(|pid| ended(pid)) {
1113 assert!(started.elapsed() < std::time::Duration::from_secs(20), "still running: {pids:?} (pty {pty})");
1114 std::thread::sleep(std::time::Duration::from_millis(20));
1115 }
1116 }
1117 }
1118
1119 #[cfg(target_os = "linux")]
1120 #[test]
1121 fn cancelling_a_child_that_shares_stdin_ends_only_the_child() {
1122 let seen = std::cell::RefCell::new(Vec::new());
1125 let outcome = Process::new("sh")
1126 .args(["-c", "sleep 60 & echo $!; wait"])
1127 .run(&|| seen.borrow().len() == 1, &mut |line| {
1128 if let Line::Out(pid) = line {
1129 seen.borrow_mut().push(pid);
1130 }
1131 })
1132 .expect("the shell starts");
1133 assert_eq!(outcome, ProcessOutcome::Cancelled);
1134 let pid = seen.into_inner().remove(0);
1135 std::thread::sleep(std::time::Duration::from_millis(200));
1136 let survived = !ended(&pid);
1137 let raw: i32 = pid.parse().expect("a process id");
1138 if let Some(pid) = rustix::process::Pid::from_raw(raw) {
1139 let _ = rustix::process::kill_process(pid, rustix::process::Signal::KILL);
1140 }
1141 assert!(survived, "the grandchild outlives a cancel of the child");
1142 }
1143
1144 #[test]
1145 fn a_line_ended_twice_by_a_return_is_kept() {
1146 let mut lines = Lines::default();
1149 let mut seen = Vec::new();
1150 lines.feed(b"hazir\r\r\nbitti\r\r\r\n", &mut |line| seen.push(line));
1151 assert_eq!(seen, vec!["hazir".to_owned(), "bitti".to_owned()]);
1152 }
1153
1154 #[cfg(unix)]
1155 #[test]
1156 fn a_pseudo_terminal_line_ended_by_the_program_itself_arrives_whole() {
1157 let (lines, _) = run(Process::new("sh").args(["-c", r"printf 'bir\r\niki\r\n'"]).pty(80, 24));
1158 assert_eq!(lines, vec![Line::Out("bir".to_owned()), Line::Out("iki".to_owned())]);
1159 }
1160
1161 #[test]
1162 fn a_line_without_an_end_is_delivered_in_pieces_of_bounded_size() {
1163 let (lines, _) = shell("head -c 300000 /dev/zero | tr '\\0' a");
1165 let total: usize = lines
1166 .iter()
1167 .map(|line| match line {
1168 Line::Out(text) => {
1169 assert!(text.len() <= MAX_LINE, "a piece of {} bytes", text.len());
1170 assert!(text.bytes().all(|byte| byte == b'a'));
1171 text.len()
1172 }
1173 Line::Err(text) => panic!("nothing was written to standard error: {text}"),
1174 })
1175 .sum();
1176 assert_eq!(total, 300_000, "nothing is lost between the pieces");
1177 }
1178
1179 #[test]
1180 fn a_long_line_is_never_cut_inside_a_character() {
1181 let mut lines = Lines::default();
1182 let mut seen = Vec::new();
1183 let mut text = vec![b'a'];
1185 for _ in 0..MAX_LINE {
1186 text.extend_from_slice("ç".as_bytes());
1187 }
1188 lines.feed(&text, &mut |line| seen.push(line));
1189 lines.finish(&mut |line| seen.push(line));
1190 assert!(seen.len() > 1, "the line was split");
1191 assert!(seen.iter().all(|line| !line.contains('\u{fffd}')), "no character was cut in two");
1192 assert_eq!(seen.concat().as_bytes(), text.as_slice());
1193 }
1194
1195 #[test]
1196 fn a_child_that_closes_its_output_can_still_be_cancelled() {
1197 let started = std::time::Instant::now();
1199 let outcome = Process::new("sh")
1200 .args(["-c", "exec >&- 2>&-; sleep 20"])
1201 .run(&|| started.elapsed() > std::time::Duration::from_millis(200), &mut |_| {})
1202 .expect("the shell starts");
1203 assert_eq!(outcome, ProcessOutcome::Cancelled);
1204 assert!(started.elapsed() < std::time::Duration::from_secs(10), "took {:?}", started.elapsed());
1205 }
1206
1207 #[test]
1208 fn a_flood_of_output_waits_for_the_reader_instead_of_piling_up() {
1209 let dir = std::env::temp_dir().join(format!("quvyta-process-flood-{}", std::process::id()));
1210 let _ = std::fs::remove_dir_all(&dir);
1211 std::fs::create_dir_all(&dir).expect("test directory");
1212 let marker = dir.join("done");
1213 let script = format!("yes | head -n 200000; touch '{}'", marker.display());
1214 let mut first = true;
1215 let mut finished_while_the_reader_slept = false;
1216 let mut count = 0_usize;
1217 let outcome = Process::new("sh")
1218 .args(["-c", &script])
1219 .run(&|| false, &mut |_| {
1220 count += 1;
1221 if first {
1222 first = false;
1223 std::thread::sleep(std::time::Duration::from_millis(700));
1225 finished_while_the_reader_slept = marker.exists();
1226 }
1227 })
1228 .expect("the shell starts");
1229 assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1230 assert_eq!(count, 200_000);
1231 assert!(!finished_while_the_reader_slept, "the child wrote everything into memory while nobody read");
1232 std::fs::remove_dir_all(&dir).expect("clean");
1233 }
1234
1235 #[cfg(target_os = "linux")]
1236 #[test]
1237 fn the_child_on_a_pseudo_terminal_holds_it_only_on_its_own_streams() {
1238 let script = r#"t=$(readlink /proc/$$/fd/1); n=0; for f in /proc/$$/fd/*; do [ "$(readlink "$f")" = "$t" ] && n=$((n+1)); done; echo $n"#;
1241 let (lines, _) = run(Process::new("sh").args(["-c", script]).pty(80, 24));
1242 assert_eq!(lines, vec![Line::Out("2".to_owned())], "standard output and standard error, nothing else");
1243 }
1244
1245 #[test]
1246 fn a_line_split_across_reads_stays_one_line() {
1247 let mut lines = Lines::default();
1248 let mut seen = Vec::new();
1249 let mut emit = |line: String| seen.push(line);
1250 lines.feed(b"ilk par", &mut emit);
1251 lines.feed(b"\xc3", &mut emit);
1252 lines.feed(b"\xa7a\r\nson", &mut emit);
1253 lines.finish(&mut emit);
1254 assert_eq!(seen, vec!["ilk parça".to_owned(), "son".to_owned()]);
1255 }
1256
1257 fn split_keeping_frames(bytes: &[u8]) -> (Vec<String>, Vec<String>) {
1260 let mut lines = Lines::default();
1261 let (mut seen, mut frames) = (Vec::new(), Vec::new());
1262 lines.feed_keeping(bytes, &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1263 lines.finish_keeping(&mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1264 (seen, frames)
1265 }
1266
1267 #[test]
1268 fn frames_overwritten_by_a_return_are_kept_only_when_asked_for() {
1269 let (seen, frames) = split_keeping_frames(b"bir\riki\ruc\rbitti\r\n");
1270 assert_eq!(seen, vec!["bitti".to_owned()]);
1271 assert_eq!(frames, vec!["bir".to_owned(), "iki".to_owned(), "uc".to_owned()]);
1272 let mut lines = Lines::default();
1273 let mut seen = Vec::new();
1274 lines.feed(b"bir\riki\ruc\rbitti\r\n", &mut |line| seen.push(line));
1275 lines.finish(&mut |line| seen.push(line));
1276 assert_eq!(seen, vec!["bitti".to_owned()]);
1277 }
1278
1279 #[test]
1280 fn a_line_ended_by_returns_and_a_newline_is_no_frame() {
1281 let (seen, frames) = split_keeping_frames(b"hazir\r\r\nbitti\r\n\rbos\r\r\r\n");
1282 assert_eq!(seen, vec!["hazir".to_owned(), "bitti".to_owned(), "bos".to_owned()]);
1283 assert!(frames.is_empty(), "{frames:?}");
1284 }
1285
1286 #[test]
1287 fn a_stream_ending_in_a_return_delivers_its_last_frame() {
1288 let (seen, frames) = split_keeping_frames(b"once\r10%\r20%\r");
1289 assert!(seen.is_empty(), "{seen:?}");
1290 assert_eq!(frames, vec!["once".to_owned(), "10%".to_owned(), "20%".to_owned()]);
1291 let mut lines = Lines::default();
1293 let (mut seen, mut frames) = (Vec::new(), Vec::new());
1294 lines.feed_keeping(b"30%\r", &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1295 assert!(frames.is_empty(), "a return before a newline is not yet known to overwrite");
1296 lines.feed_keeping(b"\n", &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1297 assert_eq!((seen, frames), (vec!["30%".to_owned()], Vec::new()));
1298 }
1299
1300 #[test]
1301 fn frames_keep_their_colour_and_erase_codes() {
1302 let (seen, frames) = split_keeping_frames(b"\x1b[1mFetch\x1b[0m 1\r\x1b[K\x1b[92mDone\x1b[0m\r\n");
1303 assert_eq!(frames, vec!["\x1b[1mFetch\x1b[0m 1".to_owned()]);
1304 assert_eq!(seen, vec!["\x1b[K\x1b[92mDone\x1b[0m".to_owned()]);
1305 }
1306
1307 #[test]
1308 fn every_frame_of_a_recorded_cargo_install_is_kept() {
1309 let recorded = include_bytes!("../../tests/fixtures/cargo-install-pty.txt");
1310 let (seen, frames) = split_keeping_frames(recorded);
1311 assert_eq!(frames.len(), 159, "every overwritten frame");
1312 assert_eq!(seen.len(), 76, "the lines themselves are unchanged");
1313 let building: Vec<&String> = frames.iter().filter(|frame| frame.contains("Building")).collect();
1314 assert_eq!(building.len(), 51);
1315 assert!(building[0].contains("] 0/46: anstyle"), "{:?}", building[0]);
1316 assert!(building.iter().any(|frame| frame.contains("] 45/46: hexyl")), "{building:?}");
1317 let mut lines = Lines::default();
1319 let mut plain = Vec::new();
1320 lines.feed(recorded, &mut |line| plain.push(line));
1321 lines.finish(&mut |line| plain.push(line));
1322 assert_eq!(plain, seen);
1323 }
1324
1325 fn run_keeping_frames(process: Process) -> Vec<(bool, Line)> {
1328 let seen = std::cell::RefCell::new(Vec::new());
1329 process
1330 .run_with_overwritten(&|| false, &mut |line| seen.borrow_mut().push((false, line)), &mut |frame| {
1331 seen.borrow_mut().push((true, frame));
1332 })
1333 .expect("the shell starts");
1334 seen.into_inner()
1335 }
1336
1337 #[test]
1338 fn overwritten_frames_arrive_through_a_pipe_in_order_and_tagged_by_stream() {
1339 let seen = run_keeping_frames(Process::new("sh").args(["-c", r"printf 'a\rb\rc\n'; printf '1%%\r2%%\r' >&2"]));
1340 let out: Vec<_> = seen.iter().filter(|(_, line)| matches!(line, Line::Out(_))).cloned().collect();
1341 let err: Vec<_> = seen.iter().filter(|(_, line)| matches!(line, Line::Err(_))).cloned().collect();
1342 assert_eq!(
1343 out,
1344 vec![
1345 (true, Line::Out("a".to_owned())),
1346 (true, Line::Out("b".to_owned())),
1347 (false, Line::Out("c".to_owned()))
1348 ]
1349 );
1350 assert_eq!(err, vec![(true, Line::Err("1%".to_owned())), (true, Line::Err("2%".to_owned()))]);
1351 }
1352
1353 #[cfg(unix)]
1354 #[test]
1355 fn overwritten_frames_arrive_from_a_pseudo_terminal() {
1356 let seen = run_keeping_frames(Process::new("sh").args(["-c", r"printf 'a\rb\rc\nd\r\n'"]).pty(80, 24));
1357 assert_eq!(
1358 seen,
1359 vec![
1360 (true, Line::Out("a".to_owned())),
1361 (true, Line::Out("b".to_owned())),
1362 (false, Line::Out("c".to_owned())),
1363 (false, Line::Out("d".to_owned())),
1364 ]
1365 );
1366 }
1367
1368 #[test]
1369 fn collected_output_arrives_in_the_order_both_streams_were_written() {
1370 let script = "printf a; printf b >&2; printf c";
1371 let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(64));
1372 assert_eq!(collected.text, "abc");
1373 assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
1374 assert!(!collected.trimmed, "nothing was dropped: {}", collected.text);
1375 assert!(!collected.timed_out && !collected.cancelled);
1376 #[cfg(unix)]
1379 {
1380 let collected = collect(Process::new("sh").args(["-c", script]).pty(80, 24), Keep::bytes(64));
1381 assert_eq!(collected.text, "abc");
1382 }
1383 }
1384
1385 #[test]
1386 fn only_the_last_bytes_of_a_flood_of_output_are_kept() {
1387 let script = "head -c 8388608 /dev/zero | tr '\\0' x; printf 'SON'";
1389 let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(65536));
1390 assert_eq!(collected.text.len(), 65536);
1391 assert!(collected.text.ends_with("SON"), "the newest bytes are the ones kept");
1392 assert!(collected.text[..65533].bytes().all(|byte| byte == b'x'), "only the flood is before them");
1393 assert!(collected.trimmed, "eight megabytes did not fit in sixty-four");
1394 }
1395
1396 #[test]
1397 fn the_lines_asked_for_are_the_last_ones() {
1398 let script = "for name in bir iki uc dort bes; do printf '%s\\n' \"$name\"; done";
1399 let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(4096).lines(3));
1400 assert_eq!(collected.text, "uc\ndort\nbes\n");
1401 assert!(collected.trimmed);
1402 let collected = collect(Process::new("sh").args(["-c", r"printf 'bir\niki\nuc'"]), Keep::bytes(4096).lines(3));
1404 assert_eq!(collected.text, "bir\niki\nuc");
1405 assert!(!collected.trimmed);
1406 let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(9).lines(3));
1408 assert_eq!(collected.text, "dort\nbes\n");
1409 }
1410
1411 #[test]
1412 fn a_tail_keeps_the_end_of_what_it_is_fed_and_never_half_a_character() {
1413 let mut tail = Tail::new(16, None);
1414 tail.feed(b"bir iki ");
1415 tail.feed(b"uc dort");
1416 assert_eq!(tail.text(), "bir iki uc dort");
1417 assert!(!tail.trimmed, "nothing fell out of sixteen bytes");
1418 tail.feed(b"!\n");
1419 assert_eq!(tail.text(), "ir iki uc dort!\n", "the oldest byte fell out");
1420 assert!(tail.trimmed);
1421 let mut tail = Tail::new(5, None);
1423 for _ in 0..8 {
1424 tail.feed("ç".as_bytes());
1425 }
1426 assert_eq!(tail.text(), "çç", "the half a character at the front is dropped");
1427 }
1428
1429 #[test]
1430 fn a_tail_keeps_the_last_lines_it_was_given() {
1431 let mut tail = Tail::new(64, Some(2));
1432 for name in ["bir\n", "iki\n", "uc\n", "dort"] {
1433 tail.feed(name.as_bytes());
1434 }
1435 assert_eq!(tail.text(), "uc\ndort", "the line being written counts as one");
1436 assert!(tail.trimmed);
1437 }
1438
1439 #[test]
1440 fn collecting_nothing_but_the_outcome_is_allowed() {
1441 let collected = collect(Process::new("sh").args(["-c", "printf 'gorunmez'"]), Keep::bytes(0));
1442 assert_eq!(collected.text, "");
1443 assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
1444 assert!(collected.trimmed);
1445 }
1446
1447 #[test]
1448 fn cancelling_a_collected_child_stops_it_the_way_a_limit_does() {
1449 let script = "sleep 30 & printf 'hazir\\n%s\\n' \"$!\"; sleep 30";
1452 let started = std::time::Instant::now();
1453 let process = Process::new("sh").args(["-c", script]).no_stdin();
1454 let collected = process
1455 .collect(Keep::bytes(1024).limit(Duration::from_secs(30)), &|| started.elapsed() > Duration::from_secs(1))
1456 .expect("the shell starts");
1457 assert!(collected.cancelled && !collected.timed_out, "{collected:?}");
1458 assert_eq!(collected.outcome, ProcessOutcome::Cancelled);
1459 assert!(started.elapsed() < Duration::from_secs(15), "took {:?}", started.elapsed());
1460 let mut lines = collected.text.lines();
1461 assert_eq!(lines.next(), Some("hazir"), "what was written before the cancel is still there: {collected:?}");
1462 #[cfg(target_os = "linux")]
1464 {
1465 let pid = lines.next().unwrap_or_default().to_owned();
1466 assert!(pid.bytes().all(|byte| byte.is_ascii_digit()), "a process id: {pid:?}");
1467 let started = std::time::Instant::now();
1468 while !ended(&pid) {
1469 assert!(started.elapsed() < Duration::from_secs(20), "the `sleep` is still running: {pid}");
1470 std::thread::sleep(Duration::from_millis(20));
1471 }
1472 }
1473 }
1474
1475 #[cfg(target_os = "linux")]
1476 #[test]
1477 fn the_limit_ends_the_child_and_the_programs_it_started() {
1478 let script = "sleep 30 & printf '%s\\n' \"$!\"; sleep 30";
1481 let started = std::time::Instant::now();
1482 let process = Process::new("sh").args(["-c", script]).no_stdin();
1483 let collected =
1484 process.collect(Keep::bytes(1024).limit(Duration::from_secs(1)), &|| false).expect("the shell starts");
1485 assert!(collected.timed_out && !collected.cancelled, "{collected:?}");
1486 assert_eq!(collected.outcome, ProcessOutcome::Finished { code: None }, "a signal ended it");
1487 assert!(started.elapsed() < Duration::from_secs(15), "took {:?}", started.elapsed());
1488 let pid = collected.text.trim().to_owned();
1490 assert!(pid.bytes().all(|byte| byte.is_ascii_digit()), "a process id: {pid:?}");
1491 let started = std::time::Instant::now();
1492 while !ended(&pid) {
1493 assert!(started.elapsed() < Duration::from_secs(20), "the `sleep` is still running: {pid}");
1494 std::thread::sleep(Duration::from_millis(20));
1495 }
1496 }
1497}