use std::collections::VecDeque;
use std::ffi::OsString;
use std::fs::File;
use std::io::{self, Read};
use std::path::PathBuf;
use std::process::{Child, Command, Stdio};
use std::sync::mpsc::{self, RecvTimeoutError, SyncSender};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::{Duration, Instant};
const POLL: Duration = Duration::from_millis(10);
pub(super) const CHUNK: usize = 4096;
const MAX_LINE: usize = 64 * 1024;
const QUEUE: usize = 1024;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Process {
program: OsString,
args: Vec<OsString>,
dir: Option<PathBuf>,
env: Vec<(OsString, OsString)>,
pty: Option<(u16, u16)>,
no_stdin: bool,
out: Option<Out>,
cleared: bool,
}
#[derive(Debug, Clone)]
struct Out(Arc<File>);
impl PartialEq for Out {
fn eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
}
impl Eq for Out {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Line {
Out(String),
Err(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ProcessOutcome {
Finished {
code: Option<i32>,
},
Cancelled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct Keep {
bytes: usize,
lines: Option<usize>,
limit: Option<Duration>,
after_exit: Option<Duration>,
}
impl Keep {
#[must_use]
pub fn bytes(bytes: usize) -> Self {
Self { bytes, lines: None, limit: None, after_exit: None }
}
#[must_use]
pub fn lines(mut self, lines: usize) -> Self {
self.lines = Some(lines);
self
}
#[must_use]
pub fn limit(mut self, limit: Duration) -> Self {
self.limit = Some(limit);
self
}
#[must_use]
pub fn after_exit(mut self, wait: Duration) -> Self {
self.after_exit = Some(wait);
self
}
}
const GRACE: Duration = Duration::from_secs(2);
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct Collected {
pub text: String,
pub outcome: ProcessOutcome,
pub trimmed: bool,
pub timed_out: bool,
pub cancelled: bool,
}
impl Process {
#[must_use]
pub fn new(program: impl Into<OsString>) -> Self {
Self {
program: program.into(),
args: Vec::new(),
dir: None,
env: Vec::new(),
pty: None,
no_stdin: false,
out: None,
cleared: false,
}
}
#[must_use]
pub fn arg(mut self, arg: impl Into<OsString>) -> Self {
self.args.push(arg.into());
self
}
#[must_use]
pub fn args(mut self, args: impl IntoIterator<Item = impl Into<OsString>>) -> Self {
self.args.extend(args.into_iter().map(Into::into));
self
}
#[must_use]
pub fn dir(mut self, dir: impl Into<PathBuf>) -> Self {
self.dir = Some(dir.into());
self
}
#[must_use]
pub fn env(mut self, key: impl Into<OsString>, value: impl Into<OsString>) -> Self {
self.env.push((key.into(), value.into()));
self
}
#[must_use]
pub fn clear_env(mut self) -> Self {
self.cleared = true;
self
}
#[must_use]
pub fn pty(mut self, cols: u16, rows: u16) -> Self {
self.pty = Some((cols, rows));
self
}
#[must_use]
pub fn no_stdin(mut self) -> Self {
self.no_stdin = true;
self
}
#[must_use]
pub fn stdout_to(mut self, file: File) -> Self {
self.out = Some(Out(Arc::new(file)));
self
}
pub fn run(self, cancel: &dyn Fn() -> bool, on_line: &mut dyn FnMut(Line)) -> io::Result<ProcessOutcome> {
self.run_inner(cancel, on_line, None)
}
pub fn run_with_overwritten(
self,
cancel: &dyn Fn() -> bool,
on_line: &mut dyn FnMut(Line),
on_overwritten: &mut dyn FnMut(Line),
) -> io::Result<ProcessOutcome> {
self.run_inner(cancel, on_line, Some(on_overwritten))
}
fn run_inner(
self,
cancel: &dyn Fn() -> bool,
on_line: &mut dyn FnMut(Line),
mut on_overwritten: Option<&mut dyn FnMut(Line)>,
) -> io::Result<ProcessOutcome> {
let frames = on_overwritten.is_some();
let group = self.no_stdin && cfg!(unix);
let command = self.command(group);
let (sender, receiver) = mpsc::sync_channel(QUEUE);
let mut child = match self.pty {
Some(size) => spawn_on_pty(command, size, &sender, frames, group)?,
None => spawn_on_pipes(command, self.out.as_ref(), &sender, frames, group)?,
};
drop(sender);
loop {
if cancel() {
kill(&mut child, group);
return Ok(ProcessOutcome::Cancelled);
}
match receiver.recv_timeout(POLL) {
Ok(Sent::Line(line)) => on_line(line),
Ok(Sent::Overwritten(frame)) => {
if let Some(on_overwritten) = on_overwritten.as_deref_mut() {
on_overwritten(frame);
}
}
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => break,
}
}
loop {
if let Some(status) = child.try_wait()? {
return Ok(ProcessOutcome::Finished { code: status.code() });
}
if cancel() {
kill(&mut child, group);
return Ok(ProcessOutcome::Cancelled);
}
std::thread::sleep(POLL);
}
}
pub fn collect(self, keep: Keep, cancel: &dyn Fn() -> bool) -> io::Result<Collected> {
self.collect_inner(keep, cancel, None)
}
pub fn collect_watching(
self,
keep: Keep,
cancel: &dyn Fn() -> bool,
mut on_tail: impl FnMut(&str),
) -> io::Result<Collected> {
self.collect_inner(keep, cancel, Some(&mut on_tail))
}
fn collect_inner(
self,
keep: Keep,
cancel: &dyn Fn() -> bool,
mut on_tail: Option<&mut dyn FnMut(&str)>,
) -> io::Result<Collected> {
let Keep { bytes, lines, limit, after_exit } = keep;
let grace = after_exit.unwrap_or(GRACE);
let group = self.no_stdin && cfg!(unix);
let command = self.command(group);
let merged = Arc::new(Mutex::new(Merged { tail: Tail::new(bytes, lines), open: true, fed: 0 }));
let (mut told, mut told_at) = (0u64, Instant::now() - WATCH_EVERY);
let mut child = match self.pty {
Some(size) => spawn_merged_on_pty(command, size, group, &merged)?,
None => spawn_merged(command, self.out.as_ref(), group, &merged)?,
};
let deadline = limit.map(|limit| Instant::now() + limit);
let mut ended: Option<(Instant, Option<i32>)> = None;
let (outcome, timed_out, cancelled) = loop {
if cancel() {
stop(&mut child, group, grace, &merged);
break (ProcessOutcome::Cancelled, false, true);
}
if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
stop(&mut child, group, grace, &merged);
break (ProcessOutcome::Finished { code: None }, true, false);
}
let open = lock(&merged).open;
if let Some(on_tail) = on_tail.as_mut()
&& told_at.elapsed() >= WATCH_EVERY
{
let fresh = {
let mut merged = lock(&merged);
(merged.fed != told).then(|| (merged.fed, merged.tail.text()))
};
if let Some((fed, text)) = fresh {
on_tail(&text);
(told, told_at) = (fed, Instant::now());
}
}
if ended.is_none()
&& let Some(status) = child.try_wait()?
{
ended = Some((Instant::now(), status.code()));
}
if let Some((at, code)) = ended {
if !open {
break (ProcessOutcome::Finished { code }, false, false);
}
if let Some(wait) = after_exit
&& at.elapsed() >= wait
{
kill(&mut child, group);
break (ProcessOutcome::Finished { code }, false, false);
}
}
std::thread::sleep(POLL);
};
let mut merged = lock(&merged);
if let Some(on_tail) = on_tail
&& merged.fed != told
{
on_tail(&merged.tail.text());
}
Ok(Collected { text: merged.tail.text(), trimmed: merged.tail.trimmed, outcome, timed_out, cancelled })
}
fn command(&self, group: bool) -> Command {
let mut command = Command::new(&self.program);
command.args(&self.args);
if self.no_stdin {
command.stdin(Stdio::null());
} else {
command.stdin(Stdio::inherit());
}
#[cfg(unix)]
if group {
use std::os::unix::process::CommandExt;
command.process_group(0);
}
#[cfg(not(unix))]
let _ = group;
if let Some(dir) = &self.dir {
command.current_dir(dir);
}
if self.cleared {
command.env_clear();
}
for (key, value) in &self.env {
command.env(key, value);
}
command
}
}
enum Sent {
Line(Line),
Overwritten(Line),
}
fn stop(child: &mut Child, group: bool, grace: Duration, merged: &Mutex<Merged>) {
#[cfg(unix)]
{
let pid = rustix::process::Pid::from_child(child);
let _ = if group {
rustix::process::kill_process_group(pid, rustix::process::Signal::TERM)
} else {
rustix::process::kill_process(pid, rustix::process::Signal::TERM)
};
let end = Instant::now() + grace;
while Instant::now() < end && (child.try_wait().ok().flatten().is_none() || lock(merged).open) {
std::thread::sleep(POLL);
}
}
#[cfg(not(unix))]
let _ = (grace, merged);
kill(child, group);
}
fn kill(child: &mut Child, group: bool) {
#[cfg(unix)]
if group {
let leader = rustix::process::Pid::from_child(child);
let _ = rustix::process::kill_process_group(leader, rustix::process::Signal::KILL);
}
#[cfg(not(unix))]
let _ = group;
let _ = child.kill();
let _ = child.wait();
}
fn spawn_on_pipes(
mut command: Command,
out: Option<&Out>,
sender: &SyncSender<Sent>,
frames: bool,
group: bool,
) -> io::Result<Child> {
match out {
Some(file) => command.stdout(Stdio::from(file.0.try_clone()?)),
None => command.stdout(Stdio::piped()),
};
command.stderr(Stdio::piped());
let mut child = command.spawn()?;
drop(command);
let started = match out {
Some(_) => match child.stderr.take() {
Some(err) => spawn_reader("err", err, Line::Err, frames, sender.clone()),
None => Err(io::Error::other("the child was started without its pipes")),
},
None => match child.stdout.take().zip(child.stderr.take()) {
Some((out, err)) => spawn_reader("out", out, Line::Out, frames, sender.clone())
.and_then(|()| spawn_reader("err", err, Line::Err, frames, sender.clone())),
None => Err(io::Error::other("the child was started without its pipes")),
},
};
match started {
Ok(()) => Ok(child),
Err(error) => {
kill(&mut child, group);
Err(error)
}
}
}
#[cfg(unix)]
fn spawn_on_pty(
mut command: Command,
size: (u16, u16),
sender: &SyncSender<Sent>,
frames: bool,
group: bool,
) -> io::Result<Child> {
let terminal = open_pty(&mut command, size)?;
let mut child = command.spawn()?;
drop(command);
match spawn_reader("pty", terminal, Line::Out, frames, sender.clone()) {
Ok(()) => Ok(child),
Err(error) => {
kill(&mut child, group);
Err(error)
}
}
}
#[cfg(not(unix))]
fn spawn_on_pty(
mut command: Command,
size: (u16, u16),
_sender: &SyncSender<Sent>,
_frames: bool,
_group: bool,
) -> io::Result<Child> {
open_pty(&mut command, size)
}
fn spawn_merged(mut command: Command, out: Option<&Out>, group: bool, merged: &Shared) -> io::Result<Child> {
let (reader, writer) = io::pipe()?;
command.stdout(match out {
Some(file) => Stdio::from(file.0.try_clone()?),
None => Stdio::from(writer.try_clone()?),
});
command.stderr(writer);
let mut child = command.spawn()?;
drop(command);
match spawn_collector(reader, Arc::clone(merged)) {
Ok(()) => Ok(child),
Err(error) => {
kill(&mut child, group);
Err(error)
}
}
}
#[cfg(unix)]
fn spawn_merged_on_pty(mut command: Command, size: (u16, u16), group: bool, merged: &Shared) -> io::Result<Child> {
let terminal = open_pty(&mut command, size)?;
let mut child = command.spawn()?;
drop(command);
match spawn_collector(terminal, Arc::clone(merged)) {
Ok(()) => Ok(child),
Err(error) => {
kill(&mut child, group);
Err(error)
}
}
}
#[cfg(not(unix))]
fn spawn_merged_on_pty(mut command: Command, size: (u16, u16), _group: bool, _merged: &Shared) -> io::Result<Child> {
open_pty(&mut command, size)
}
#[cfg(unix)]
fn open_pty(command: &mut Command, (cols, rows): (u16, u16)) -> io::Result<std::fs::File> {
use std::os::fd::OwnedFd;
use rustix::fs::{Mode, OFlags};
use rustix::io::{FdFlags, fcntl_setfd};
use rustix::pty::{OpenptFlags, grantpt, openpt, ptsname, unlockpt};
use rustix::termios::{Winsize, tcsetwinsize};
#[cfg(any(target_os = "linux", target_os = "android", target_os = "freebsd", target_os = "netbsd"))]
let flags = OpenptFlags::RDWR | OpenptFlags::NOCTTY | OpenptFlags::CLOEXEC;
#[cfg(not(any(target_os = "linux", target_os = "android", target_os = "freebsd", target_os = "netbsd")))]
let flags = OpenptFlags::RDWR | OpenptFlags::NOCTTY;
let controller = openpt(flags)?;
fcntl_setfd(&controller, FdFlags::CLOEXEC)?;
grantpt(&controller)?;
unlockpt(&controller)?;
tcsetwinsize(&controller, Winsize { ws_row: rows, ws_col: cols, ws_xpixel: 0, ws_ypixel: 0 })?;
let name = ptsname(&controller, Vec::new())?;
let device: OwnedFd = rustix::fs::open(name, OFlags::RDWR | OFlags::NOCTTY | OFlags::CLOEXEC, Mode::empty())?;
command.stdout(Stdio::from(device.try_clone()?)).stderr(Stdio::from(device));
Ok(File::from(controller))
}
#[cfg(not(unix))]
fn open_pty(_command: &mut Command, _size: (u16, u16)) -> io::Result<std::fs::File> {
Err(io::Error::new(io::ErrorKind::Unsupported, "a pseudo-terminal needs a Unix system"))
}
fn spawn_reader(
name: &str,
source: impl Read + Send + 'static,
tag: fn(String) -> Line,
frames: bool,
sender: SyncSender<Sent>,
) -> io::Result<()> {
std::thread::Builder::new()
.name(format!("quvyta-process-{name}"))
.spawn(move || read_lines(source, tag, frames, &sender))
.map(|_| ())
}
fn read_lines(mut source: impl Read, tag: fn(String) -> Line, frames: bool, sender: &SyncSender<Sent>) {
let mut chunk = [0_u8; CHUNK];
let mut lines = Lines::default();
let listening = std::cell::Cell::new(true);
let mut on_line = |line| listening.set(listening.get() && sender.send(Sent::Line(tag(line))).is_ok());
let mut on_frame = |frame| listening.set(listening.get() && sender.send(Sent::Overwritten(tag(frame))).is_ok());
loop {
match source.read(&mut chunk) {
Ok(0) => break,
Ok(count) => {
lines.feed_keeping(&chunk[..count], &mut on_line, frames.then_some(&mut on_frame));
if !listening.get() {
return;
}
}
Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
Err(_) => break,
}
}
lines.finish_keeping(&mut on_line, frames.then_some(&mut on_frame));
}
fn spawn_collector(source: impl Read + Send + 'static, merged: Shared) -> io::Result<()> {
std::thread::Builder::new()
.name("quvyta-process-merged".to_owned())
.spawn(move || read_tail(source, &merged))
.map(|_| ())
}
#[derive(Debug)]
struct Merged {
tail: Tail,
open: bool,
fed: u64,
}
const WATCH_EVERY: Duration = Duration::from_millis(100);
type Shared = Arc<Mutex<Merged>>;
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
fn read_tail(mut source: impl Read, merged: &Shared) {
let mut chunk = [0_u8; CHUNK];
loop {
match source.read(&mut chunk) {
Ok(0) => break,
Ok(count) => {
let mut merged = lock(merged);
merged.tail.feed(&chunk[..count]);
merged.fed += 1;
}
Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
Err(_) => break,
}
}
lock(merged).open = false;
}
#[derive(Debug)]
struct Tail {
bytes: VecDeque<u8>,
lines: usize,
ends_line: bool,
kept: usize,
lines_kept: Option<usize>,
trimmed: bool,
}
impl Tail {
fn new(bytes: usize, lines: Option<usize>) -> Self {
Self { bytes: VecDeque::new(), lines: 0, ends_line: true, kept: bytes, lines_kept: lines, trimmed: false }
}
fn feed(&mut self, bytes: &[u8]) {
let Some(&last) = bytes.last() else {
return;
};
self.ends_line = last == b'\n';
self.lines += bytes.iter().filter(|&&byte| byte == b'\n').count();
self.bytes.extend(bytes);
if let Some(kept) = self.lines_kept {
for _ in 0..self.counting().saturating_sub(kept) {
self.drop_oldest_line();
}
}
for _ in 0..self.bytes.len().saturating_sub(self.kept) {
if self.bytes.pop_front() == Some(b'\n') {
self.lines -= 1;
}
self.trimmed = true;
}
}
fn counting(&self) -> usize {
self.lines + usize::from(!self.ends_line && !self.bytes.is_empty())
}
fn drop_oldest_line(&mut self) {
self.trimmed = true;
while let Some(byte) = self.bytes.pop_front() {
if byte == b'\n' {
self.lines -= 1;
break;
}
}
}
fn text(&mut self) -> String {
let bytes = self.bytes.make_contiguous();
let start = bytes.iter().copied().take(3).take_while(|byte| byte & 0b1100_0000 == 0b1000_0000).count();
String::from_utf8_lossy(&bytes[start..]).into_owned()
}
}
#[derive(Debug, Default)]
pub(super) struct Lines {
buffer: Vec<u8>,
pending_return: bool,
}
impl Lines {
pub(super) fn feed(&mut self, bytes: &[u8], emit: &mut impl FnMut(String)) {
self.feed_keeping(bytes, emit, None);
}
pub(super) fn feed_keeping(
&mut self,
bytes: &[u8],
emit: &mut impl FnMut(String),
mut overwritten: Option<&mut dyn FnMut(String)>,
) {
for &byte in bytes {
if self.pending_return {
match byte {
b'\r' => continue,
b'\n' => {
self.pending_return = false;
emit(self.take());
continue;
}
_ => {
self.pending_return = false;
self.overwrite(&mut overwritten);
}
}
}
match byte {
b'\r' => self.pending_return = true,
b'\n' => emit(self.take()),
_ => {
self.buffer.push(byte);
if self.buffer.len() >= MAX_LINE {
self.emit_piece(emit);
}
}
}
}
}
fn emit_piece(&mut self, emit: &mut impl FnMut(String)) {
let len = self.buffer.len();
let mut cut = len;
for back in 1..=len.min(3) {
let byte = self.buffer[len - back];
if byte & 0b1100_0000 != 0b1000_0000 {
let width = match byte {
0xc0..=0xdf => 2,
0xe0..=0xef => 3,
0xf0..=0xf7 => 4,
_ => 1,
};
if width > back {
cut = len - back;
}
break;
}
}
let rest = self.buffer.split_off(cut);
emit(self.take());
self.buffer = rest;
}
fn overwrite(&mut self, overwritten: &mut Option<&mut dyn FnMut(String)>) {
match overwritten {
Some(overwritten) if !self.buffer.is_empty() => overwritten(self.take()),
_ => self.buffer.clear(),
}
}
pub(super) fn finish(&mut self, emit: &mut impl FnMut(String)) {
self.finish_keeping(emit, None);
}
pub(super) fn finish_keeping(
&mut self,
emit: &mut impl FnMut(String),
mut overwritten: Option<&mut dyn FnMut(String)>,
) {
if self.pending_return {
self.overwrite(&mut overwritten);
self.pending_return = false;
}
if !self.buffer.is_empty() {
emit(self.take());
}
}
fn take(&mut self) -> String {
let line = String::from_utf8_lossy(&self.buffer).into_owned();
self.buffer.clear();
line
}
}
#[cfg(test)]
mod tests {
use std::fs::File;
use std::path::PathBuf;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use super::{Collected, Keep, Line, Lines, MAX_LINE, Process, ProcessOutcome, Tail};
fn shell(script: &str) -> (Vec<Line>, ProcessOutcome) {
run(Process::new("sh").args(["-c", script]))
}
fn run(process: Process) -> (Vec<Line>, ProcessOutcome) {
let mut lines = Vec::new();
let outcome = process.run(&|| false, &mut |line| lines.push(line)).expect("the shell starts");
(lines, outcome)
}
fn collect(process: Process, keep: Keep) -> Collected {
process.collect(keep, &|| false).expect("the shell starts")
}
fn folder(name: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!("quvyta-process-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("test folder");
dir
}
#[test]
fn keeps_the_two_streams_apart_and_reports_the_exit_code() {
let (lines, outcome) = shell("echo bir; echo iki >&2; exit 3");
assert_eq!(lines.len(), 2, "{lines:?}");
assert!(lines.contains(&Line::Out("bir".to_owned())), "{lines:?}");
assert!(lines.contains(&Line::Err("iki".to_owned())), "{lines:?}");
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(3) });
}
#[test]
fn delivers_the_last_line_without_a_newline() {
let (lines, outcome) = shell("printf 'son satir'");
assert_eq!(lines, vec![Line::Out("son satir".to_owned())]);
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
}
#[test]
fn carriage_returns_collapse_into_one_line() {
let (lines, _) = shell(r"printf 'a\rbb\rccc\n'");
assert_eq!(lines, vec![Line::Out("ccc".to_owned())]);
}
#[test]
fn invalid_utf8_becomes_the_replacement_character() {
let (lines, _) = shell(r"printf 'a\377b\n'");
assert_eq!(lines, vec![Line::Out("a\u{fffd}b".to_owned())]);
}
#[test]
fn the_environment_is_inherited_and_one_variable_can_be_replaced() {
let (lines, _) = shell("echo ${PATH:+inherited}");
assert_eq!(lines, vec![Line::Out("inherited".to_owned())]);
let (lines, _) = run(Process::new("sh").args(["-c", "echo $LC_ALL"]).env("LC_ALL", "C"));
assert_eq!(lines, vec![Line::Out("C".to_owned())]);
}
#[test]
fn a_cleared_environment_gives_the_child_only_what_it_was_told() {
let script = r#"printf '%s' "${HOME-unset}""#;
let (lines, _) = shell(script);
assert_ne!(lines, vec![Line::Out("unset".to_owned())], "the test process really has a home");
let path = std::env::var("PATH").expect("the test process was started with a path");
let (lines, _) = run(Process::new("sh").args(["-c", script]).clear_env().env("PATH", path.clone()));
assert_eq!(lines, vec![Line::Out("unset".to_owned())], "nothing is inherited");
let (lines, _) = run(Process::new("sh")
.args(["-c", r#"printf '%s' "${PATH-unset}""#])
.clear_env()
.env("PATH", path.clone()));
assert_eq!(lines, vec![Line::Out(path)]);
let (lines, _) =
run(Process::new("sh").args(["-c", r#"printf '%s' "${EV-unset}""#]).clear_env().env("EV", "1"));
assert_eq!(lines, vec![Line::Out("1".to_owned())], "a variable that was given arrives");
}
#[test]
fn runs_in_the_directory_it_is_given() {
let (lines, _) = run(Process::new("sh").args(["-c", "pwd"]).dir("/"));
assert_eq!(lines, vec![Line::Out("/".to_owned())]);
}
#[test]
fn cancelling_kills_a_long_running_child() {
let seen = AtomicUsize::new(0);
let outcome = Process::new("sh")
.args(["-c", "while true; do echo tik; sleep 0.05; done"])
.run(&|| seen.load(Ordering::Relaxed) > 0, &mut |line| {
assert_eq!(line, Line::Out("tik".to_owned()));
seen.fetch_add(1, Ordering::Relaxed);
})
.expect("the shell starts");
assert_eq!(outcome, ProcessOutcome::Cancelled);
assert!(seen.load(Ordering::Relaxed) > 0);
}
#[test]
fn a_missing_program_is_an_error_and_not_a_panic() {
let error = Process::new("quvyta-no-such-program")
.run(&|| false, &mut |_| unreachable!("a missing program writes nothing"))
.expect_err("a missing program cannot run");
assert_eq!(error.kind(), std::io::ErrorKind::NotFound);
}
#[cfg(unix)]
#[test]
fn on_a_pseudo_terminal_the_child_sees_a_terminal_of_the_size_we_gave() {
let (lines, outcome) = run(Process::new("sh").args(["-c", "test -t 1 && stty size <&1"]).pty(100, 24));
assert_eq!(lines, vec![Line::Out("24 100".to_owned())]);
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
}
#[cfg(unix)]
#[test]
fn on_a_pseudo_terminal_both_streams_arrive_as_output() {
let (lines, outcome) = run(Process::new("sh").args(["-c", "echo bir; echo iki >&2"]).pty(80, 24));
assert_eq!(lines, vec![Line::Out("bir".to_owned()), Line::Out("iki".to_owned())]);
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
}
#[cfg(target_os = "linux")]
fn stat_ids(stat: &str) -> [String; 3] {
let fields: Vec<&str> = stat[stat.rfind(')').expect("name") + 2..].split(' ').collect();
[fields[2], fields[3], fields[4]].map(str::to_owned)
}
#[cfg(target_os = "linux")]
#[test]
fn without_stdin_the_child_reads_an_empty_stream() {
let script = r#"readlink /proc/$$/fd/0; read answer; echo "read $?""#;
for process in [Process::new("sh").args(["-c", script]), Process::new("sh").args(["-c", script]).pty(80, 24)] {
let (lines, outcome) = run(process.no_stdin());
assert_eq!(lines, vec![Line::Out("/dev/null".to_owned()), Line::Out("read 1".to_owned())]);
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
}
}
#[cfg(target_os = "linux")]
#[test]
fn only_a_child_without_stdin_gets_a_group_of_its_own_and_it_keeps_the_session() {
let script = "cat /proc/$$/stat";
let ours = stat_ids(&std::fs::read_to_string("/proc/self/stat").expect("stat"));
let ids = |process: Process| {
let (lines, _) = run(process);
let [Line::Out(stat)] = &lines[..] else { panic!("one line: {lines:?}") };
stat_ids(stat)
};
let shared = ids(Process::new("sh").args(["-c", script]));
assert_eq!(shared, ours, "a child reading the terminal stays in the application's group");
for process in [Process::new("sh").args(["-c", script]), Process::new("sh").args(["-c", script]).pty(80, 24)] {
let [group, session, terminal] = ids(process.no_stdin());
assert_ne!(group, ours[0], "a group of its own");
assert_eq!(session, ours[1], "the application's session");
assert_eq!(terminal, ours[2], "the application's controlling terminal");
}
}
#[cfg(target_os = "linux")]
fn ended(pid: &str) -> bool {
std::fs::read_to_string(format!("/proc/{pid}/stat"))
.map_or(true, |stat| stat[stat.rfind(')').expect("name") + 2..].starts_with('Z'))
}
#[cfg(target_os = "linux")]
#[test]
fn cancelling_a_child_without_stdin_ends_the_programs_it_started() {
for pty in [false, true] {
let seen = std::cell::RefCell::new(Vec::new());
let process = Process::new("sh").args(["-c", "sleep 60 & echo $!; sleep 60 & echo $!; wait"]).no_stdin();
let process = if pty { process.pty(80, 24) } else { process };
let outcome = process
.run(&|| seen.borrow().len() == 2, &mut |line| match line {
Line::Out(pid) => seen.borrow_mut().push(pid),
Line::Err(text) => panic!("nothing on standard error: {text}"),
})
.expect("the shell starts");
assert_eq!(outcome, ProcessOutcome::Cancelled);
let pids = seen.into_inner();
let started = std::time::Instant::now();
while !pids.iter().all(|pid| ended(pid)) {
assert!(started.elapsed() < std::time::Duration::from_secs(20), "still running: {pids:?} (pty {pty})");
std::thread::sleep(std::time::Duration::from_millis(20));
}
}
}
#[cfg(target_os = "linux")]
#[test]
fn cancelling_a_child_that_shares_stdin_ends_only_the_child() {
let seen = std::cell::RefCell::new(Vec::new());
let outcome = Process::new("sh")
.args(["-c", "sleep 60 & echo $!; wait"])
.run(&|| seen.borrow().len() == 1, &mut |line| {
if let Line::Out(pid) = line {
seen.borrow_mut().push(pid);
}
})
.expect("the shell starts");
assert_eq!(outcome, ProcessOutcome::Cancelled);
let pid = seen.into_inner().remove(0);
std::thread::sleep(std::time::Duration::from_millis(200));
let survived = !ended(&pid);
let raw: i32 = pid.parse().expect("a process id");
if let Some(pid) = rustix::process::Pid::from_raw(raw) {
let _ = rustix::process::kill_process(pid, rustix::process::Signal::KILL);
}
assert!(survived, "the grandchild outlives a cancel of the child");
}
#[test]
fn a_line_ended_twice_by_a_return_is_kept() {
let mut lines = Lines::default();
let mut seen = Vec::new();
lines.feed(b"hazir\r\r\nbitti\r\r\r\n", &mut |line| seen.push(line));
assert_eq!(seen, vec!["hazir".to_owned(), "bitti".to_owned()]);
}
#[cfg(unix)]
#[test]
fn a_pseudo_terminal_line_ended_by_the_program_itself_arrives_whole() {
let (lines, _) = run(Process::new("sh").args(["-c", r"printf 'bir\r\niki\r\n'"]).pty(80, 24));
assert_eq!(lines, vec![Line::Out("bir".to_owned()), Line::Out("iki".to_owned())]);
}
#[test]
fn a_line_without_an_end_is_delivered_in_pieces_of_bounded_size() {
let (lines, _) = shell("head -c 300000 /dev/zero | tr '\\0' a");
let total: usize = lines
.iter()
.map(|line| match line {
Line::Out(text) => {
assert!(text.len() <= MAX_LINE, "a piece of {} bytes", text.len());
assert!(text.bytes().all(|byte| byte == b'a'));
text.len()
}
Line::Err(text) => panic!("nothing was written to standard error: {text}"),
})
.sum();
assert_eq!(total, 300_000, "nothing is lost between the pieces");
}
#[test]
fn a_long_line_is_never_cut_inside_a_character() {
let mut lines = Lines::default();
let mut seen = Vec::new();
let mut text = vec![b'a'];
for _ in 0..MAX_LINE {
text.extend_from_slice("ç".as_bytes());
}
lines.feed(&text, &mut |line| seen.push(line));
lines.finish(&mut |line| seen.push(line));
assert!(seen.len() > 1, "the line was split");
assert!(seen.iter().all(|line| !line.contains('\u{fffd}')), "no character was cut in two");
assert_eq!(seen.concat().as_bytes(), text.as_slice());
}
#[test]
fn a_child_that_closes_its_output_can_still_be_cancelled() {
let started = std::time::Instant::now();
let outcome = Process::new("sh")
.args(["-c", "exec >&- 2>&-; sleep 20"])
.run(&|| started.elapsed() > std::time::Duration::from_millis(200), &mut |_| {})
.expect("the shell starts");
assert_eq!(outcome, ProcessOutcome::Cancelled);
assert!(started.elapsed() < std::time::Duration::from_secs(10), "took {:?}", started.elapsed());
}
#[test]
fn a_flood_of_output_waits_for_the_reader_instead_of_piling_up() {
let dir = folder("flood");
let marker = dir.join("done");
let script = format!("yes | head -n 200000; touch '{}'", marker.display());
let mut first = true;
let mut finished_while_the_reader_slept = false;
let mut count = 0_usize;
let outcome = Process::new("sh")
.args(["-c", &script])
.run(&|| false, &mut |_| {
count += 1;
if first {
first = false;
std::thread::sleep(std::time::Duration::from_millis(700));
finished_while_the_reader_slept = marker.exists();
}
})
.expect("the shell starts");
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
assert_eq!(count, 200_000);
assert!(!finished_while_the_reader_slept, "the child wrote everything into memory while nobody read");
std::fs::remove_dir_all(&dir).expect("clean");
}
#[cfg(target_os = "linux")]
#[test]
fn the_child_on_a_pseudo_terminal_holds_it_only_on_its_own_streams() {
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"#;
let (lines, _) = run(Process::new("sh").args(["-c", script]).pty(80, 24));
assert_eq!(lines, vec![Line::Out("2".to_owned())], "standard output and standard error, nothing else");
}
#[test]
fn a_line_split_across_reads_stays_one_line() {
let mut lines = Lines::default();
let mut seen = Vec::new();
let mut emit = |line: String| seen.push(line);
lines.feed(b"ilk par", &mut emit);
lines.feed(b"\xc3", &mut emit);
lines.feed(b"\xa7a\r\nson", &mut emit);
lines.finish(&mut emit);
assert_eq!(seen, vec!["ilk parça".to_owned(), "son".to_owned()]);
}
fn split_keeping_frames(bytes: &[u8]) -> (Vec<String>, Vec<String>) {
let mut lines = Lines::default();
let (mut seen, mut frames) = (Vec::new(), Vec::new());
lines.feed_keeping(bytes, &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
lines.finish_keeping(&mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
(seen, frames)
}
#[test]
fn frames_overwritten_by_a_return_are_kept_only_when_asked_for() {
let (seen, frames) = split_keeping_frames(b"bir\riki\ruc\rbitti\r\n");
assert_eq!(seen, vec!["bitti".to_owned()]);
assert_eq!(frames, vec!["bir".to_owned(), "iki".to_owned(), "uc".to_owned()]);
let mut lines = Lines::default();
let mut seen = Vec::new();
lines.feed(b"bir\riki\ruc\rbitti\r\n", &mut |line| seen.push(line));
lines.finish(&mut |line| seen.push(line));
assert_eq!(seen, vec!["bitti".to_owned()]);
}
#[test]
fn a_line_ended_by_returns_and_a_newline_is_no_frame() {
let (seen, frames) = split_keeping_frames(b"hazir\r\r\nbitti\r\n\rbos\r\r\r\n");
assert_eq!(seen, vec!["hazir".to_owned(), "bitti".to_owned(), "bos".to_owned()]);
assert!(frames.is_empty(), "{frames:?}");
}
#[test]
fn a_stream_ending_in_a_return_delivers_its_last_frame() {
let (seen, frames) = split_keeping_frames(b"once\r10%\r20%\r");
assert!(seen.is_empty(), "{seen:?}");
assert_eq!(frames, vec!["once".to_owned(), "10%".to_owned(), "20%".to_owned()]);
let mut lines = Lines::default();
let (mut seen, mut frames) = (Vec::new(), Vec::new());
lines.feed_keeping(b"30%\r", &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
assert!(frames.is_empty(), "a return before a newline is not yet known to overwrite");
lines.feed_keeping(b"\n", &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
assert_eq!((seen, frames), (vec!["30%".to_owned()], Vec::new()));
}
#[test]
fn frames_keep_their_colour_and_erase_codes() {
let (seen, frames) = split_keeping_frames(b"\x1b[1mFetch\x1b[0m 1\r\x1b[K\x1b[92mDone\x1b[0m\r\n");
assert_eq!(frames, vec!["\x1b[1mFetch\x1b[0m 1".to_owned()]);
assert_eq!(seen, vec!["\x1b[K\x1b[92mDone\x1b[0m".to_owned()]);
}
#[test]
fn every_frame_of_a_recorded_cargo_install_is_kept() {
let recorded = include_bytes!("../../tests/fixtures/cargo-install-pty.txt");
let (seen, frames) = split_keeping_frames(recorded);
assert_eq!(frames.len(), 159, "every overwritten frame");
assert_eq!(seen.len(), 76, "the lines themselves are unchanged");
let building: Vec<&String> = frames.iter().filter(|frame| frame.contains("Building")).collect();
assert_eq!(building.len(), 51);
assert!(building[0].contains("] 0/46: anstyle"), "{:?}", building[0]);
assert!(building.iter().any(|frame| frame.contains("] 45/46: hexyl")), "{building:?}");
let mut lines = Lines::default();
let mut plain = Vec::new();
lines.feed(recorded, &mut |line| plain.push(line));
lines.finish(&mut |line| plain.push(line));
assert_eq!(plain, seen);
}
fn run_keeping_frames(process: Process) -> Vec<(bool, Line)> {
let seen = std::cell::RefCell::new(Vec::new());
process
.run_with_overwritten(&|| false, &mut |line| seen.borrow_mut().push((false, line)), &mut |frame| {
seen.borrow_mut().push((true, frame));
})
.expect("the shell starts");
seen.into_inner()
}
#[test]
fn overwritten_frames_arrive_through_a_pipe_in_order_and_tagged_by_stream() {
let seen = run_keeping_frames(Process::new("sh").args(["-c", r"printf 'a\rb\rc\n'; printf '1%%\r2%%\r' >&2"]));
let out: Vec<_> = seen.iter().filter(|(_, line)| matches!(line, Line::Out(_))).cloned().collect();
let err: Vec<_> = seen.iter().filter(|(_, line)| matches!(line, Line::Err(_))).cloned().collect();
assert_eq!(
out,
vec![
(true, Line::Out("a".to_owned())),
(true, Line::Out("b".to_owned())),
(false, Line::Out("c".to_owned()))
]
);
assert_eq!(err, vec![(true, Line::Err("1%".to_owned())), (true, Line::Err("2%".to_owned()))]);
}
#[cfg(unix)]
#[test]
fn overwritten_frames_arrive_from_a_pseudo_terminal() {
let seen = run_keeping_frames(Process::new("sh").args(["-c", r"printf 'a\rb\rc\nd\r\n'"]).pty(80, 24));
assert_eq!(
seen,
vec![
(true, Line::Out("a".to_owned())),
(true, Line::Out("b".to_owned())),
(false, Line::Out("c".to_owned())),
(false, Line::Out("d".to_owned())),
]
);
}
#[test]
fn collected_output_arrives_in_the_order_both_streams_were_written() {
let script = "printf a; printf b >&2; printf c";
let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(64));
assert_eq!(collected.text, "abc");
assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
assert!(!collected.trimmed, "nothing was dropped: {}", collected.text);
assert!(!collected.timed_out && !collected.cancelled);
#[cfg(unix)]
{
let collected = collect(Process::new("sh").args(["-c", script]).pty(80, 24), Keep::bytes(64));
assert_eq!(collected.text, "abc");
}
}
#[test]
fn only_the_last_bytes_of_a_flood_of_output_are_kept() {
let script = "head -c 8388608 /dev/zero | tr '\\0' x; printf 'SON'";
let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(65536));
assert_eq!(collected.text.len(), 65536);
assert!(collected.text.ends_with("SON"), "the newest bytes are the ones kept");
assert!(collected.text[..65533].bytes().all(|byte| byte == b'x'), "only the flood is before them");
assert!(collected.trimmed, "eight megabytes did not fit in sixty-four");
}
#[test]
fn the_lines_asked_for_are_the_last_ones() {
let script = "for name in bir iki uc dort bes; do printf '%s\\n' \"$name\"; done";
let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(4096).lines(3));
assert_eq!(collected.text, "uc\ndort\nbes\n");
assert!(collected.trimmed);
let collected = collect(Process::new("sh").args(["-c", r"printf 'bir\niki\nuc'"]), Keep::bytes(4096).lines(3));
assert_eq!(collected.text, "bir\niki\nuc");
assert!(!collected.trimmed);
let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(9).lines(3));
assert_eq!(collected.text, "dort\nbes\n");
}
#[test]
fn a_tail_keeps_the_end_of_what_it_is_fed_and_never_half_a_character() {
let mut tail = Tail::new(16, None);
tail.feed(b"bir iki ");
tail.feed(b"uc dort");
assert_eq!(tail.text(), "bir iki uc dort");
assert!(!tail.trimmed, "nothing fell out of sixteen bytes");
tail.feed(b"!\n");
assert_eq!(tail.text(), "ir iki uc dort!\n", "the oldest byte fell out");
assert!(tail.trimmed);
let mut tail = Tail::new(5, None);
for _ in 0..8 {
tail.feed("ç".as_bytes());
}
assert_eq!(tail.text(), "çç", "the half a character at the front is dropped");
}
#[test]
fn a_tail_keeps_the_last_lines_it_was_given() {
let mut tail = Tail::new(64, Some(2));
for name in ["bir\n", "iki\n", "uc\n", "dort"] {
tail.feed(name.as_bytes());
}
assert_eq!(tail.text(), "uc\ndort", "the line being written counts as one");
assert!(tail.trimmed);
}
#[test]
fn collecting_nothing_but_the_outcome_is_allowed() {
let collected = collect(Process::new("sh").args(["-c", "printf 'gorunmez'"]), Keep::bytes(0));
assert_eq!(collected.text, "");
assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
assert!(collected.trimmed);
}
#[test]
fn the_standard_output_can_go_straight_to_a_file_and_only_the_error_stream_is_kept() {
let dir = folder("into-file");
let path = dir.join("unpacked.bin");
let script = "head -c 8388608 /dev/zero; \
n=0; while [ \"$n\" -lt 300 ]; do printf 'satir %s\\n' \"$n\" >&2; n=$((n+1)); done";
let file = File::create(&path).expect("the file is created");
let collected = collect(Process::new("sh").args(["-c", script]).stdout_to(file), Keep::bytes(64));
assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
assert_eq!(
std::fs::metadata(&path).expect("the file is there").len(),
8_388_608,
"every byte of it is in the file"
);
assert_eq!(collected.text.len(), 64, "{}", collected.text);
assert!(collected.text.ends_with("satir 299\n"), "the newest line is kept: {}", collected.text);
assert!(!collected.text.contains("satir 1"), "the older lines fell out: {}", collected.text);
assert!(collected.trimmed);
std::fs::remove_dir_all(&dir).expect("clean");
}
#[test]
fn a_file_named_for_the_output_leaves_the_error_stream_the_only_one_to_read() {
let dir = folder("into-file-lines");
let path = dir.join("out.txt");
let file = File::create(&path).expect("the file is created");
let (lines, outcome) = run(Process::new("sh").args(["-c", "printf cikti; printf hata >&2"]).stdout_to(file));
assert_eq!(lines, vec![Line::Err("hata".to_owned())], "the output went to the file: {lines:?}");
assert_eq!(std::fs::read_to_string(&path).expect("the file is there"), "cikti");
assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
std::fs::remove_dir_all(&dir).expect("clean");
}
#[test]
fn a_cancelled_child_writes_nothing_more_to_the_file() {
let dir = folder("into-file-cancel");
let path = dir.join("flood.bin");
let file = File::create(&path).expect("the file is created");
let script = "while :; do printf '0123456789012345678901234567890123456789'; done";
let started = std::time::Instant::now();
let process = Process::new("sh").args(["-c", script]).no_stdin().stdout_to(file);
let collected = process
.collect(Keep::bytes(1024), &|| started.elapsed() > Duration::from_millis(500))
.expect("the shell starts");
assert!(collected.cancelled && !collected.timed_out, "{collected:?}");
let size = || std::fs::metadata(&path).expect("the file is there").len();
let stopped = size();
assert!(stopped > 0, "the child wrote into the file before it was cancelled");
for _ in 0..10 {
std::thread::sleep(Duration::from_millis(200));
assert_eq!(size(), stopped, "the file grew after the child was cancelled");
}
std::fs::remove_dir_all(&dir).expect("clean");
}
#[test]
fn a_command_that_leaves_something_behind_is_done_once_it_has_ended() {
let dir = folder("after-exit");
let started = std::time::Instant::now();
let collected = Process::new("sh")
.args(["-c", "(sleep 4; printf late > late.txt) & printf done"])
.dir(&dir)
.no_stdin()
.collect(Keep::bytes(1024).after_exit(Duration::from_secs(2)), &|| false)
.expect("the shell starts");
assert!(started.elapsed() < Duration::from_secs(15), "it did not wait for what it left behind");
assert_eq!(collected.text, "done");
assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
assert!(!collected.timed_out && !collected.cancelled, "{collected:?}");
std::thread::sleep(Duration::from_secs(6));
assert!(!dir.join("late.txt").exists(), "the program left behind was ended");
std::fs::remove_dir_all(&dir).expect("clean");
}
#[test]
fn a_stopped_child_is_asked_first_and_has_its_last_words_kept() {
let started = std::time::Instant::now();
let collected = Process::new("sh")
.args(["-c", "trap 'printf bye; exit 0' TERM; sleep 30 & wait"])
.no_stdin()
.collect(Keep::bytes(1024), &|| started.elapsed() > Duration::from_secs(1))
.expect("the shell starts");
assert!(collected.cancelled, "{collected:?}");
assert!(collected.text.contains("bye"), "TERM came first and what it said is kept: {collected:?}");
assert!(started.elapsed() < Duration::from_secs(15), "and the stop is bounded");
}
#[test]
fn a_watcher_sees_the_output_go_by_while_the_child_runs() {
let mut seen = Vec::new();
let collected = Process::new("sh")
.args(["-c", "printf a; sleep 1; printf b; sleep 1; printf c"])
.no_stdin()
.collect_watching(Keep::bytes(1024), &|| false, |tail| seen.push(tail.to_owned()))
.expect("the shell starts");
assert_eq!(collected.text, "abc");
assert!(seen.len() >= 2, "told more than once while it ran: {seen:?}");
assert!(seen.iter().any(|tail| tail.contains('a') && !tail.contains('c')), "before the end: {seen:?}");
assert_eq!(seen.last().map(String::as_str), Some("abc"), "and the last call has it all: {seen:?}");
}
#[test]
fn a_watcher_sees_no_more_lines_than_are_kept_and_nothing_of_a_silent_child() {
let mut seen = Vec::new();
Process::new("sh")
.args(["-c", "printf '1\\n2\\n'; sleep 0.3; printf '3\\n4\\n'"])
.no_stdin()
.collect_watching(Keep::bytes(1024).lines(2), &|| false, |tail| seen.push(tail.to_owned()))
.expect("the shell starts");
assert!(!seen.is_empty());
assert!(seen.iter().all(|tail| tail.lines().count() <= 2), "{seen:?}");
let mut silent = 0;
Process::new("sleep")
.args(["1"])
.no_stdin()
.collect_watching(Keep::bytes(1024), &|| false, |_| silent += 1)
.expect("sleep starts");
assert_eq!(silent, 0, "nothing written, nothing told");
}
#[test]
fn cancelling_a_collected_child_stops_it_the_way_a_limit_does() {
let script = "sleep 30 & printf 'hazir\\n%s\\n' \"$!\"; sleep 30";
let started = std::time::Instant::now();
let process = Process::new("sh").args(["-c", script]).no_stdin();
let collected = process
.collect(Keep::bytes(1024).limit(Duration::from_secs(30)), &|| started.elapsed() > Duration::from_secs(1))
.expect("the shell starts");
assert!(collected.cancelled && !collected.timed_out, "{collected:?}");
assert_eq!(collected.outcome, ProcessOutcome::Cancelled);
assert!(started.elapsed() < Duration::from_secs(15), "took {:?}", started.elapsed());
let mut lines = collected.text.lines();
assert_eq!(lines.next(), Some("hazir"), "what was written before the cancel is still there: {collected:?}");
#[cfg(target_os = "linux")]
{
let pid = lines.next().unwrap_or_default().to_owned();
assert!(pid.bytes().all(|byte| byte.is_ascii_digit()), "a process id: {pid:?}");
let started = std::time::Instant::now();
while !ended(&pid) {
assert!(started.elapsed() < Duration::from_secs(20), "the `sleep` is still running: {pid}");
std::thread::sleep(Duration::from_millis(20));
}
}
}
#[cfg(target_os = "linux")]
#[test]
fn the_limit_ends_the_child_and_the_programs_it_started() {
let script = "sleep 30 & printf '%s\\n' \"$!\"; sleep 30";
let started = std::time::Instant::now();
let process = Process::new("sh").args(["-c", script]).no_stdin();
let collected =
process.collect(Keep::bytes(1024).limit(Duration::from_secs(1)), &|| false).expect("the shell starts");
assert!(collected.timed_out && !collected.cancelled, "{collected:?}");
assert_eq!(collected.outcome, ProcessOutcome::Finished { code: None }, "a signal ended it");
assert!(started.elapsed() < Duration::from_secs(15), "took {:?}", started.elapsed());
let pid = collected.text.trim().to_owned();
assert!(pid.bytes().all(|byte| byte.is_ascii_digit()), "a process id: {pid:?}");
let started = std::time::Instant::now();
while !ended(&pid) {
assert!(started.elapsed() < Duration::from_secs(20), "the `sleep` is still running: {pid}");
std::thread::sleep(Duration::from_millis(20));
}
}
}