#[cfg(any(not(target_os = "linux"), test))]
use std::{
io::{self, Read, Write},
process::{Command, Output, Stdio},
thread,
time::{Duration, Instant},
};
#[cfg(any(not(target_os = "linux"), test))]
pub fn output_within(mut command: Command, input: &[u8], limit: Duration) -> io::Result<Output> {
let mut child = command
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()?;
fn drain<R: Read + Send + 'static>(pipe: Option<R>) -> thread::JoinHandle<Vec<u8>> {
thread::spawn(move || {
let mut bytes = Vec::new();
if let Some(mut pipe) = pipe {
let _ = pipe.read_to_end(&mut bytes);
}
bytes
})
}
let stdout = drain(child.stdout.take());
let stderr = drain(child.stderr.take());
if let Some(mut stdin) = child.stdin.take() {
match stdin.write_all(input) {
Err(e) if e.kind() != io::ErrorKind::BrokenPipe => return Err(e),
_ => {}
}
}
let deadline = Instant::now() + limit;
let status = loop {
if let Some(status) = child.try_wait()? {
break status;
}
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
return Err(io::Error::new(
io::ErrorKind::TimedOut,
format!("no answer within {}s", limit.as_secs()),
));
}
thread::sleep(Duration::from_millis(5));
};
Ok(Output {
status,
stdout: stdout.join().unwrap_or_default(),
stderr: stderr.join().unwrap_or_default(),
})
}
#[cfg(target_os = "linux")]
pub fn running(program: &str) -> Option<usize> {
let mut names = String::new();
for entry in std::fs::read_dir("/proc").ok()?.flatten() {
if entry
.file_name()
.to_string_lossy()
.bytes()
.all(|b| b.is_ascii_digit())
&& let Ok(comm) = std::fs::read_to_string(entry.path().join("comm"))
{
names.push_str(comm.trim());
names.push('\n');
}
}
Some(count_named(&names, program))
}
#[cfg(not(target_os = "linux"))]
pub fn running(program: &str) -> Option<usize> {
let mut ps = Command::new("/bin/ps");
ps.args(["-A", "-o", "comm="]);
let out = output_within(ps, b"", Duration::from_secs(5)).ok()?;
if !out.status.success() {
return None;
}
Some(count_named(&String::from_utf8_lossy(&out.stdout), program))
}
fn count_named(listing: &str, program: &str) -> usize {
listing
.lines()
.map(str::trim)
.filter(|line| {
std::path::Path::new(line)
.file_name()
.is_some_and(|name| name == program)
})
.count()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_program_is_counted_by_its_own_name_and_nothing_else() {
let listing = "/Users/a/.codex/packages/standalone/bin/codex\n\
codex\n\
/bin/zsh\n\
/usr/bin/codex-helper\n\
node\n";
assert_eq!(count_named(listing, "codex"), 2);
assert_eq!(count_named(listing, "gemini"), 0);
}
#[test]
fn the_process_list_can_be_read() {
assert_eq!(running("definitely-not-a-program-name"), Some(0));
let me = std::env::current_exe().expect("this test's own binary");
let name = me.file_name().unwrap().to_string_lossy().into_owned();
let name: String = if cfg!(target_os = "linux") {
name.chars().take(15).collect()
} else {
name
};
assert!(running(&name).is_some_and(|n| n >= 1), "{name} is running");
}
#[test]
fn a_prompt_answer_is_returned_whole() {
let mut cat = Command::new("cat");
cat.arg("-");
let out = output_within(cat, b"hello", Duration::from_secs(5)).unwrap();
assert!(out.status.success());
assert_eq!(out.stdout, b"hello");
}
#[test]
fn a_helper_that_never_answers_is_stopped_at_the_deadline() {
let mut sleep = Command::new("sleep");
sleep.arg("30");
let started = Instant::now();
let err = output_within(sleep, b"", Duration::from_millis(200)).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::TimedOut);
assert!(
started.elapsed() < Duration::from_secs(5),
"{:?}",
started.elapsed()
);
}
#[test]
fn a_large_answer_cannot_stall_the_helper() {
let mut yes = Command::new("head");
yes.args(["-c", "1000000", "/dev/zero"]);
let out = output_within(yes, b"", Duration::from_secs(10)).unwrap();
assert_eq!(out.stdout.len(), 1_000_000);
}
}