use super::ExecError;
use std::ffi::OsString;
use std::io::{Read, Write};
use std::path::Path;
use std::process::{Child, Command, ExitStatus, Stdio};
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::{Duration, Instant};
pub(super) const ETXTBSY_RETRY_ATTEMPTS: u32 = 100;
const POLL_INTERVAL: Duration = Duration::from_millis(50);
const ETXTBSY_RETRY_INTERVAL: Duration = Duration::from_millis(2);
pub(super) struct Captured {
pub(super) stdout: Vec<u8>,
pub(super) stderr: Vec<u8>,
pub(super) status: ExitStatus,
}
pub(super) struct SpawnArgs<'a> {
pub(super) binary: &'a OsString,
pub(super) args: &'a [OsString],
pub(super) stdin_bytes: &'a [u8],
pub(super) extra_env: &'a [(&'a str, OsString)],
pub(super) cwd: &'a Path,
pub(super) stop: &'a AtomicBool,
pub(super) deadline: Duration,
pub(super) etxtbsy_budget: u32,
pub(super) tool_name: &'a str,
}
pub(super) fn spawn_and_capture(req: &SpawnArgs<'_>) -> Result<Captured, ExecError> {
let mut child = spawn_with_etxtbsy_retry(req)?;
let stdin_data = req.stdin_bytes.to_vec();
let mut child_stdin = child.stdin.take().expect("stdin is piped");
let stdin_thread = thread::spawn(move || {
let _ = child_stdin.write_all(&stdin_data);
drop(child_stdin);
});
let mut child_stdout = child.stdout.take().expect("stdout is piped");
let mut child_stderr = child.stderr.take().expect("stderr is piped");
let stdout_thread = thread::spawn(move || {
let mut buf = Vec::new();
let _ = child_stdout.read_to_end(&mut buf);
buf
});
let stderr_thread = thread::spawn(move || {
let mut buf = Vec::new();
let _ = child_stderr.read_to_end(&mut buf);
buf
});
let status = wait_with_stop(&mut child, req.stop, req.deadline);
stdin_thread.join().expect("stdin writer did not panic");
let stdout = stdout_thread.join().expect("stdout reader did not panic");
let stderr = stderr_thread.join().expect("stderr reader did not panic");
Ok(Captured {
stdout,
stderr,
status,
})
}
fn spawn_with_etxtbsy_retry(req: &SpawnArgs<'_>) -> Result<Child, ExecError> {
let mut attempt: u32 = 1;
loop {
let mut cmd = Command::new(req.binary);
cmd.args(req.args)
.current_dir(req.cwd)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
for (k, v) in req.extra_env {
cmd.env(k, v);
}
match cmd.spawn() {
Ok(child) => return Ok(child),
Err(e) if e.raw_os_error() == Some(libc::ETXTBSY) && attempt < req.etxtbsy_budget => {
attempt += 1;
thread::sleep(ETXTBSY_RETRY_INTERVAL);
}
Err(source) => {
return Err(ExecError::Spawn {
name: req.tool_name.to_string(),
source,
});
}
}
}
}
fn wait_with_stop(child: &mut Child, stop: &AtomicBool, deadline: Duration) -> ExitStatus {
loop {
if let Some(status) = try_reap(child) {
return status;
}
thread::sleep(POLL_INTERVAL);
if stop.load(Ordering::SeqCst) {
return cascade_terminate(child, deadline);
}
}
}
fn try_reap(child: &mut Child) -> Option<ExitStatus> {
child.try_wait().ok().flatten()
}
fn cascade_terminate(child: &mut Child, deadline: Duration) -> ExitStatus {
let pid = child.id() as i32;
unsafe {
libc::kill(pid, libc::SIGTERM);
}
let term_until = Instant::now() + deadline;
while Instant::now() < term_until {
if let Some(status) = try_reap(child) {
return status;
}
thread::sleep(POLL_INTERVAL);
}
unsafe {
libc::kill(pid, libc::SIGKILL);
}
child.wait().expect("kernel reaps SIGKILL'd child")
}