Files
nucleic-brush/brush-core/src/processes.rs
T

130 lines
4.5 KiB
Rust
Raw Normal View History

//! Process management
use futures::FutureExt;
use crate::{error, sys};
/// A waitable future that will yield the results of a child process's execution.
pub(crate) type WaitableChildProcess = std::pin::Pin<
Box<dyn futures::Future<Output = Result<std::process::Output, std::io::Error>> + Send + Sync>,
>;
/// Tracks a child process being awaited.
pub struct ChildProcess {
/// A waitable future that will yield the results of a child process's execution.
exec_future: WaitableChildProcess,
/// If available, the process ID of the child.
pid: Option<sys::process::ProcessId>,
/// If available, the process group ID of the child.
pgid: Option<sys::process::ProcessId>,
// nash: gate token issued by `Gate::on_exec` for this command, reported
// back via `Gate::on_exit` when the process completes.
pub(crate) gate_token: Option<u64>,
}
impl ChildProcess {
/// Wraps a child process and its future.
pub fn new(
child: sys::process::Child,
pid: Option<sys::process::ProcessId>,
pgid: Option<sys::process::ProcessId>,
) -> Self {
Self {
exec_future: Box::pin(child.wait_with_output()),
pid,
pgid,
gate_token: None,
}
}
/// Returns the process's ID.
pub const fn pid(&self) -> Option<sys::process::ProcessId> {
self.pid
}
/// Returns the process's group ID.
pub const fn pgid(&self) -> Option<sys::process::ProcessId> {
self.pgid
}
/// Waits for the process to exit.
pub async fn wait(&mut self) -> Result<ProcessWaitResult, error::Error> {
#[allow(unused_mut, reason = "only mutated on some platforms")]
let mut sigtstp = sys::signal::tstp_signal_listener()?;
#[allow(unused_mut, reason = "only mutated on some platforms")]
let mut sigchld = sys::signal::chld_signal_listener()?;
#[allow(clippy::ignored_unit_patterns)]
loop {
tokio::select! {
output = &mut self.exec_future => {
let output = output?;
// nash: report process completion to the gate exactly once.
if let Some(token) = self.gate_token.take() {
crate::gate::gate().on_exit(
token,
crate::gate::ExecEnd { exit_code: exit_code_of(&output.status) },
);
}
break Ok(ProcessWaitResult::Completed(output))
},
_ = sigtstp.recv() => {
break Ok(ProcessWaitResult::Stopped)
},
_ = sigchld.recv() => {
if sys::signal::poll_for_stopped_children()? {
break Ok(ProcessWaitResult::Stopped);
}
},
_ = sys::signal::await_ctrl_c() => {
// SIGINT got thrown. Handle it and continue looping. The child should
// have received it as well, and either handled it or ended up getting
// terminated (in which case we'll see the child exit).
},
}
}
}
pub(crate) fn poll(&mut self) -> Option<Result<std::process::Output, error::Error>> {
let checkable_future = &mut self.exec_future;
let result: Option<Result<std::process::Output, error::Error>> = checkable_future
.now_or_never()
.map(|result| result.map_err(Into::into));
// nash: completion can also be observed via poll (e.g. job checks).
if let Some(Ok(output)) = &result {
if let Some(token) = self.gate_token.take() {
crate::gate::gate().on_exit(
token,
crate::gate::ExecEnd {
exit_code: exit_code_of(&output.status),
},
);
}
}
result
}
}
// nash: numeric exit code for gate reporting (128+signal for signal deaths).
fn exit_code_of(status: &std::process::ExitStatus) -> i32 {
if let Some(code) = status.code() {
return code;
}
#[cfg(unix)]
{
use std::os::unix::process::ExitStatusExt;
if let Some(signal) = status.signal() {
return 128 + signal;
}
}
1
}
/// Represents the result of waiting for an executing process.
pub enum ProcessWaitResult {
/// The process completed.
Completed(std::process::Output),
/// The process stopped and has not yet completed.
Stopped,
}