use std::io::Read;
use std::ops::{Deref, DerefMut};
use std::path::Path;
use std::process::{Child, Command, Stdio};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::{Duration, Instant};
pub type SharedBuf = Arc<Mutex<String>>;
pub fn drain_to_buffer<R: Read + Send + 'static>(mut reader: R, buffer: SharedBuf) {
let mut chunk = [0u8; 4096];
loop {
match reader.read(&mut chunk) {
Ok(0) => return, Ok(n) => {
let text = String::from_utf8_lossy(&chunk[..n]);
let mut guard = buffer.lock().expect("buffer lock poisoned");
guard.push_str(&text);
}
Err(e) => {
eprintln!("reader thread io error: {e}");
return;
}
}
}
}
#[allow(dead_code)]
pub struct Drained {
out_handle: thread::JoinHandle<()>,
err_handle: thread::JoinHandle<()>,
out_buf: SharedBuf,
err_buf: SharedBuf,
}
#[allow(dead_code)]
impl Drained {
pub fn captured(&self) -> String {
format!(
"stdout:\n{}\nstderr:\n{}",
self.out_buf.lock().expect("stdout buffer lock poisoned"),
self.err_buf.lock().expect("stderr buffer lock poisoned")
)
}
pub fn finish(self) -> String {
let Drained {
out_handle,
err_handle,
out_buf,
err_buf,
} = self;
let _ = out_handle.join();
let _ = err_handle.join();
format!(
"stdout:\n{}\nstderr:\n{}",
out_buf.lock().expect("stdout buffer lock poisoned"),
err_buf.lock().expect("stderr buffer lock poisoned")
)
}
pub fn markers(&self) -> [SharedBuf; 2] {
[Arc::clone(&self.out_buf), Arc::clone(&self.err_buf)]
}
}
#[allow(dead_code)]
pub fn spawn_drained(child: &mut Child) -> Drained {
let out_buf: SharedBuf = Arc::new(Mutex::new(String::new()));
let err_buf: SharedBuf = Arc::new(Mutex::new(String::new()));
let stdout = child
.stdout
.take()
.expect("child stdout was configured as piped");
let stderr = child
.stderr
.take()
.expect("child stderr was configured as piped");
let out_handle = thread::spawn({
let buf = Arc::clone(&out_buf);
move || drain_to_buffer(stdout, buf)
});
let err_handle = thread::spawn({
let buf = Arc::clone(&err_buf);
move || drain_to_buffer(stderr, buf)
});
Drained {
out_handle,
err_handle,
out_buf,
err_buf,
}
}
pub struct KillOnDrop(pub Child);
impl Drop for KillOnDrop {
fn drop(&mut self) {
let _ = self.0.kill();
let _ = self.0.wait();
}
}
impl Deref for KillOnDrop {
type Target = Child;
fn deref(&self) -> &Child {
&self.0
}
}
impl DerefMut for KillOnDrop {
fn deref_mut(&mut self) -> &mut Child {
&mut self.0
}
}
#[allow(dead_code)]
pub fn wait_for_marker(
child: &mut Child,
buffers: &[SharedBuf],
marker: &str,
timeout: Duration,
) -> bool {
let start = Instant::now();
let step = Duration::from_millis(25);
loop {
if buffers
.iter()
.any(|buf| buf.lock().expect("buffer lock poisoned").contains(marker))
{
return true;
}
if start.elapsed() >= timeout {
return false;
}
if let Ok(Some(_)) = child.try_wait() {
return false;
}
thread::sleep(step);
}
}
#[allow(dead_code)]
pub fn send_term(child: &Child) {
let status = Command::new("kill")
.arg("-TERM")
.arg(child.id().to_string())
.status()
.expect("failed to spawn `kill -TERM`");
assert!(
status.success(),
"`kill -TERM` returned non-zero: {status:?}"
);
}
#[allow(dead_code)]
pub fn send_signal(child: &Child, signal: &str) {
let status = Command::new("kill")
.arg(signal)
.arg(child.id().to_string())
.status()
.expect("failed to spawn `kill`");
assert!(
status.success(),
"`kill {signal}` returned non-zero: {status:?}"
);
}
#[allow(dead_code)]
pub fn spawn_camel_run(dir: &Path) -> KillOnDrop {
let config_path = dir.join("Camel.toml");
let child = Command::new(env!("CARGO_BIN_EXE_camel"))
.arg("run")
.arg("--config")
.arg(&config_path)
.current_dir(dir)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.stdin(Stdio::null())
.spawn()
.expect("failed to spawn `camel` binary"); KillOnDrop(child)
}
#[allow(dead_code)]
pub fn spawn_camel_job(dir: &Path, doc: &Path) -> KillOnDrop {
spawn_camel_job_with_args(dir, doc, &[])
}
#[allow(dead_code)]
pub fn spawn_camel_job_with_args(dir: &Path, doc: &Path, extra_args: &[&str]) -> KillOnDrop {
let mut command = Command::new(env!("CARGO_BIN_EXE_camel"));
command.arg("job").arg(doc);
for arg in extra_args {
command.arg(arg);
}
let child = command
.current_dir(dir)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.stdin(Stdio::null())
.spawn()
.expect("failed to spawn `camel` binary"); KillOnDrop(child)
}
#[allow(dead_code)]
fn wait_exit_core(child: &mut Child, timeout: Duration) -> Option<i32> {
let start = Instant::now();
let step = Duration::from_millis(25);
loop {
match child.try_wait() {
Ok(Some(status)) => return Some(status.code().unwrap_or(-1)),
Ok(None) => {
if start.elapsed() >= timeout {
let _ = child.kill();
let _ = child.wait();
return None;
}
thread::sleep(step);
}
Err(e) => panic!("try_wait failed: {e}"),
}
}
}
#[allow(dead_code)]
pub fn wait_exit_bounded(child: &mut Child, timeout: Duration) -> bool {
wait_exit_core(child, timeout).is_some()
}
#[allow(dead_code)]
pub fn wait_exit_code_bounded(child: &mut Child, timeout: Duration) -> i32 {
wait_exit_core(child, timeout).unwrap_or(-1)
}
#[allow(dead_code)]
pub fn run_binary(
dir: &Path,
program: &Path,
args: &[&str],
envs: &[(&str, &str)],
) -> (i32, String, String) {
let mut child = Command::new(program)
.args(args)
.envs(envs.iter().copied())
.current_dir(dir)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.stdin(Stdio::null())
.spawn()
.expect("spawn child process");
let out_buf = SharedBuf::default();
let err_buf = SharedBuf::default();
let out_handle = thread::spawn({
let buf = Arc::clone(&out_buf);
let stdout = child.stdout.take().expect("stdout piped");
move || drain_to_buffer(stdout, buf)
});
let err_handle = thread::spawn({
let buf = Arc::clone(&err_buf);
let stderr = child.stderr.take().expect("stderr piped");
move || drain_to_buffer(stderr, buf)
});
let deadline = Instant::now() + Duration::from_secs(90);
let exit_code = loop {
match child.try_wait() {
Ok(Some(status)) => break status.code().unwrap_or(-1),
Ok(None) => {
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
break -1;
}
thread::sleep(Duration::from_millis(25));
}
Err(e) => panic!("try_wait failed: {e}"),
}
};
let _ = out_handle.join();
let _ = err_handle.join();
(
exit_code,
out_buf.lock().expect("stdout lock").clone(),
err_buf.lock().expect("stderr lock").clone(),
)
}