use std::os::unix::io::RawFd;
use std::time::{Duration, Instant};
use crate::CoreError;
use crate::error::syscall_ret;
use crate::fd::Fd;
use crate::io::DrainState;
use crate::reactor::Reactor;
use libc::{
O_CLOEXEC, O_NONBLOCK, WEXITSTATUS, WIFEXITED, WIFSIGNALED, WTERMSIG, pid_t, pipe2, waitpid,
};
mod exec;
mod fork;
mod posix;
use exec::ExecContext;
use fork::spawn_fork_internal;
use posix::spawn_posix_internal;
unsafe extern "C" {
pub(crate) static mut environ: *mut *mut libc::c_char;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum CancelPolicy {
#[default]
None,
Graceful,
Kill,
}
#[derive(Debug, Clone, Copy, Default)]
pub struct ProcessGroup {
pub leader: Option<pid_t>,
pub isolated: bool,
}
impl ProcessGroup {
pub fn new(leader: Option<pid_t>, isolated: bool) -> Self {
Self { leader, isolated }
}
}
#[inline(always)]
fn errno() -> i32 {
std::io::Error::last_os_error().raw_os_error().unwrap_or(0)
}
#[inline(always)]
fn make_pipe() -> Result<(Fd, Fd), CoreError> {
let mut fds = [0; 2];
let r = unsafe { pipe2(fds.as_mut_ptr(), O_CLOEXEC | O_NONBLOCK) };
syscall_ret(r, "pipe2")?;
Ok((Fd::new(fds[0], "pipe2")?, Fd::new(fds[1], "pipe2")?))
}
fn make_cloexec_pipe() -> Result<(RawFd, RawFd), CoreError> {
let mut fds = [0; 2];
let r = unsafe { pipe2(fds.as_mut_ptr(), O_CLOEXEC) };
syscall_ret(r, "pipe2")?;
Ok((fds[0], fds[1]))
}
struct Pipes {
stdin_r: Option<Fd>,
stdin_w: Option<Fd>,
stdout_r: Option<Fd>,
stdout_w: Option<Fd>,
stderr_r: Option<Fd>,
stderr_w: Option<Fd>,
}
impl Pipes {
fn new(in_buf: Option<&[u8]>, out: bool, err: bool) -> Result<Self, CoreError> {
let (stdin_r, stdin_w) = if in_buf.is_some() {
let (r, w) = make_pipe()?;
(Some(r), Some(w))
} else {
(None, None)
};
let (stdout_r, stdout_w) = if out {
let (r, w) = make_pipe()?;
(Some(r), Some(w))
} else {
(None, None)
};
let (stderr_r, stderr_w) = if err {
let (r, w) = make_pipe()?;
(Some(r), Some(w))
} else {
(None, None)
};
Ok(Self {
stdin_r,
stdin_w,
stdout_r,
stdout_w,
stderr_r,
stderr_w,
})
}
#[inline(always)]
fn close_all(&mut self) {
self.stdin_r.take();
self.stdin_w.take();
self.stdout_r.take();
self.stdout_w.take();
self.stderr_r.take();
self.stderr_w.take();
}
}
#[derive(Debug, PartialEq, Eq)]
pub enum ExitStatus {
Exited(i32),
Signaled(i32),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SpawnBackend {
PosixSpawn,
Fork,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum SpawnFdPolicy {
#[default]
CloexecOnly,
CloseFrom3,
Allowlist(Vec<RawFd>),
}
#[inline(always)]
fn decode_status(status: i32) -> ExitStatus {
if WIFEXITED(status) {
ExitStatus::Exited(WEXITSTATUS(status))
} else if WIFSIGNALED(status) {
ExitStatus::Signaled(WTERMSIG(status))
} else {
ExitStatus::Exited(-1)
}
}
pub struct Process {
pid: pid_t,
}
impl Process {
pub fn new(pid: pid_t) -> Self {
Self { pid }
}
pub fn pid(&self) -> pid_t {
self.pid
}
pub fn wait_step(&self) -> Result<Option<ExitStatus>, CoreError> {
loop {
let mut status = 0;
let r = unsafe { waitpid(self.pid, &mut status, libc::WNOHANG) };
if r == 0 {
return Ok(None);
}
if r < 0 {
let e = errno();
if e == libc::EINTR {
continue;
}
return Err(CoreError::sys(e, "waitpid_step"));
}
return Ok(Some(decode_status(status)));
}
}
pub fn wait_blocking(&self) -> Result<ExitStatus, CoreError> {
loop {
let mut status = 0;
let r = unsafe { waitpid(self.pid, &mut status, 0) };
if r < 0 {
let e = errno();
if e == libc::EINTR {
continue;
}
return Err(CoreError::sys(e, "waitpid_blocking"));
}
return Ok(decode_status(status));
}
}
pub fn kill(&self, sig: i32) -> Result<(), CoreError> {
let r = unsafe { libc::kill(self.pid, sig) };
if r < 0 {
let e = errno();
if e == libc::ESRCH {
return Ok(());
}
syscall_ret(-1, "kill")?;
}
Ok(())
}
pub fn kill_pgroup(&self, sig: i32) -> Result<(), CoreError> {
self.kill_group(self.pid, sig)
}
pub fn kill_group(&self, pgid: pid_t, sig: i32) -> Result<(), CoreError> {
let r = unsafe { libc::kill(-pgid, sig) };
if r < 0 {
let e = errno();
if e == libc::ESRCH {
return Ok(());
}
syscall_ret(-1, "kill_group")?;
}
Ok(())
}
}
#[derive(Clone)]
pub struct SpawnOptions {
ctx: ExecContext,
stdin: Option<Box<[u8]>>,
capture_stdout: bool,
capture_stderr: bool,
wait: bool,
pgroup: ProcessGroup,
max_output: usize,
timeout_ms: Option<u32>,
kill_grace_ms: u32,
cancel: CancelPolicy,
backend: SpawnBackend,
fd_policy: SpawnFdPolicy,
early_exit: Option<fn(&[u8]) -> bool>,
}
impl SpawnOptions {
pub fn builder(argv: Vec<String>, backend: SpawnBackend) -> SpawnOptionsBuilder {
SpawnOptionsBuilder::new(argv, backend)
}
pub fn run(self) -> Result<Output, CoreError> {
spawn(self)
}
}
#[derive(Clone)]
pub struct SpawnOptionsBuilder {
argv: Vec<String>,
env: Option<Vec<String>>,
cwd: Option<String>,
stdin: Option<Box<[u8]>>,
capture_stdout: bool,
capture_stderr: bool,
wait: bool,
pgroup: ProcessGroup,
max_output: usize,
timeout_ms: Option<u32>,
kill_grace_ms: u32,
cancel: CancelPolicy,
backend: SpawnBackend,
fd_policy: SpawnFdPolicy,
early_exit: Option<fn(&[u8]) -> bool>,
}
impl SpawnOptionsBuilder {
pub fn new(argv: Vec<String>, backend: SpawnBackend) -> Self {
Self {
argv,
env: None,
cwd: None,
stdin: None,
capture_stdout: false,
capture_stderr: false,
wait: true,
pgroup: ProcessGroup::default(),
max_output: 1024 * 1024,
timeout_ms: None,
kill_grace_ms: 2000,
cancel: CancelPolicy::Kill,
backend,
fd_policy: SpawnFdPolicy::default(),
early_exit: None,
}
}
pub fn env(mut self, env: Vec<String>) -> Self {
self.env = Some(env);
self
}
pub fn cwd(mut self, cwd: String) -> Self {
self.cwd = Some(cwd);
self
}
pub fn stdin(mut self, data: impl Into<Box<[u8]>>) -> Self {
self.stdin = Some(data.into());
self
}
pub fn capture_stdout(mut self) -> Self {
self.capture_stdout = true;
self
}
pub fn capture_stderr(mut self) -> Self {
self.capture_stderr = true;
self
}
pub fn wait(mut self, wait: bool) -> Self {
self.wait = wait;
self
}
pub fn pgroup(mut self, pgroup: ProcessGroup) -> Self {
self.pgroup = pgroup;
self
}
pub fn max_output(mut self, max: usize) -> Self {
self.max_output = max;
self
}
pub fn timeout_ms(mut self, ms: u32) -> Self {
self.timeout_ms = Some(ms);
self
}
pub fn kill_grace_ms(mut self, ms: u32) -> Self {
self.kill_grace_ms = ms;
self
}
pub fn cancel(mut self, policy: CancelPolicy) -> Self {
self.cancel = policy;
self
}
pub fn fd_policy(mut self, policy: SpawnFdPolicy) -> Self {
self.fd_policy = policy;
self
}
pub fn early_exit(mut self, callback: fn(&[u8]) -> bool) -> Self {
self.early_exit = Some(callback);
self
}
pub fn build(self) -> Result<SpawnOptions, CoreError> {
let ctx = ExecContext::new(self.argv, self.env, self.cwd)?;
Ok(SpawnOptions {
ctx,
stdin: self.stdin,
capture_stdout: self.capture_stdout,
capture_stderr: self.capture_stderr,
wait: self.wait,
pgroup: self.pgroup,
max_output: self.max_output,
timeout_ms: self.timeout_ms,
kill_grace_ms: self.kill_grace_ms,
cancel: self.cancel,
backend: self.backend,
fd_policy: self.fd_policy,
early_exit: self.early_exit,
})
}
}
#[derive(Debug)]
pub struct Output {
pub pid: pid_t,
pub status: Option<ExitStatus>,
pub stdout: Vec<u8>,
pub stderr: Vec<u8>,
pub timed_out: bool,
pub stdout_early_exited: bool,
}
fn validate_backend(opts: &SpawnOptions) -> Result<(), CoreError> {
validate_fd_policy(&opts.fd_policy)?;
match opts.backend {
SpawnBackend::PosixSpawn => {
if opts.ctx.cwd.is_some() {
return Err(CoreError::sys(libc::EINVAL, "posix_spawn cwd unsupported"));
}
if opts.pgroup.isolated {
return Err(CoreError::sys(
libc::EINVAL,
"posix_spawn setsid unsupported",
));
}
if opts.fd_policy != SpawnFdPolicy::CloexecOnly {
return Err(CoreError::sys(
libc::EINVAL,
"posix_spawn fd policy unsupported",
));
}
Ok(())
}
SpawnBackend::Fork => Ok(()),
}
}
fn validate_fd_policy(policy: &SpawnFdPolicy) -> Result<(), CoreError> {
if let SpawnFdPolicy::Allowlist(fds) = policy {
let mut seen = Vec::with_capacity(fds.len());
for &fd in fds {
if fd < 0 {
return Err(CoreError::sys(libc::EINVAL, "spawn fd allowlist invalid"));
}
let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
if flags < 0 {
return Err(CoreError::sys(errno(), "spawn fd allowlist fcntl(F_GETFD)"));
}
if seen.contains(&fd) {
return Err(CoreError::sys(libc::EINVAL, "spawn fd allowlist duplicate"));
}
seen.push(fd);
}
}
Ok(())
}
pub type SpawnDrain = DrainState<fn(&[u8]) -> bool>;
pub struct RunningProcess {
pub process: Process,
drain: SpawnDrain,
}
pub struct ManagedProcess {
running: Option<RunningProcess>,
timeout_at: Option<Instant>,
kill_grace: Duration,
cancel: CancelPolicy,
pgroup: ProcessGroup,
cancel_at: Option<Instant>,
kill_state: KillState,
status: Option<ExitStatus>,
timed_out: bool,
}
impl RunningProcess {
pub fn register_with_reactor(&mut self, reactor: &mut Reactor) -> Result<(), CoreError> {
self.drain.register_with_reactor(reactor)
}
pub fn handle_reactor_event(
&mut self,
reactor: &mut Reactor,
event: &crate::fd::Event,
) -> Result<(), CoreError> {
if self.drain.stdout_matches(event.token) {
if event.readable || event.hangup {
self.drain.handle_stdout_ready(reactor)?;
} else if event.error {
self.drain.drop_stdout(reactor)?;
}
} else if self.drain.stderr_matches(event.token) {
if event.readable || event.hangup {
self.drain.handle_stderr_ready(reactor)?;
} else if event.error {
self.drain.drop_stderr(reactor)?;
}
} else if self.drain.stdin_matches(event.token) {
if event.writable {
self.drain.handle_stdin_writable(reactor)?;
} else if event.error || event.hangup {
self.drain.drop_stdin(reactor)?;
}
}
Ok(())
}
pub fn io_done(&self) -> bool {
self.drain.is_done()
}
pub fn into_output_parts(self) -> (Vec<u8>, Vec<u8>) {
self.drain.into_parts()
}
}
impl ManagedProcess {
pub fn pid(&self) -> pid_t {
self.running
.as_ref()
.expect("managed process already completed")
.process
.pid()
}
pub fn register_with_reactor(&mut self, reactor: &mut Reactor) -> Result<(), CoreError> {
self.running
.as_mut()
.ok_or_else(|| CoreError::sys(libc::EINVAL, "managed process completed"))?
.register_with_reactor(reactor)
}
pub fn handle_reactor_event(
&mut self,
reactor: &mut Reactor,
event: &crate::fd::Event,
) -> Result<(), CoreError> {
self.running
.as_mut()
.ok_or_else(|| CoreError::sys(libc::EINVAL, "managed process completed"))?
.handle_reactor_event(reactor, event)
}
pub fn request_cancel(&mut self) {
self.cancel_at.get_or_insert_with(Instant::now);
}
pub fn next_deadline(&self) -> Option<Instant> {
self.running.as_ref()?;
let now = Instant::now();
let mut next = now + Duration::from_millis(100);
if !self.timed_out
&& let Some(timeout_at) = self.timeout_at
&& timeout_at < next
{
next = timeout_at;
}
if self.kill_state == KillState::TermSent
&& let Some(cancel_at) = self.cancel_at
{
let kill_at = cancel_at + self.kill_grace;
if kill_at < next {
next = kill_at;
}
}
Some(next)
}
pub fn poll_completion(&mut self, reactor: &mut Reactor) -> Result<Option<Output>, CoreError> {
let now = Instant::now();
if !self.timed_out
&& let Some(timeout_at) = self.timeout_at
&& now >= timeout_at
{
self.timed_out = true;
self.cancel_at.get_or_insert(timeout_at);
}
self.advance_cancel(now)?;
let running = self
.running
.as_ref()
.ok_or_else(|| CoreError::sys(libc::EINVAL, "managed process completed"))?;
if self.status.is_none() {
self.status = running.process.wait_step()?;
}
let io_done = running.io_done();
if self.status.is_some() && (io_done || self.cancel_at.is_some()) {
return self.finish(reactor, !io_done).map(Some);
}
Ok(None)
}
fn advance_cancel(&mut self, now: Instant) -> Result<(), CoreError> {
let Some(cancel_at) = self.cancel_at else {
return Ok(());
};
let running = self
.running
.as_ref()
.ok_or_else(|| CoreError::sys(libc::EINVAL, "managed process completed"))?;
let process = &running.process;
let pid = process.pid();
let pgid = effective_pgid(pid, self.pgroup);
let target_is_group = self.pgroup.isolated || self.pgroup.leader.is_some();
match self.kill_state {
KillState::None => match self.cancel {
CancelPolicy::None => {}
CancelPolicy::Graceful => {
let result = signal_process(process, target_is_group, pgid, libc::SIGTERM);
self.kill_state = if result.is_ok() {
KillState::TermSent
} else {
KillState::KillSent
};
}
CancelPolicy::Kill => {
let _ = signal_process(process, target_is_group, pgid, libc::SIGKILL);
self.kill_state = KillState::KillSent;
}
},
KillState::TermSent if now >= cancel_at + self.kill_grace => {
let _ = signal_process(process, target_is_group, pgid, libc::SIGKILL);
self.kill_state = KillState::KillSent;
}
_ => {}
}
Ok(())
}
fn finish(&mut self, reactor: &mut Reactor, force_close: bool) -> Result<Output, CoreError> {
let mut running = self
.running
.take()
.ok_or_else(|| CoreError::sys(libc::EINVAL, "managed process completed"))?;
for slot in running.drain.take_all_slots() {
if force_close {
let _ = reactor.del(&slot.fd);
} else {
reactor.del(&slot.fd)?;
}
}
let pid = running.process.pid();
let (stdout, stderr, output_limit_exceeded, stdout_early_exited) =
running.drain.into_parts_with_state();
if output_limit_exceeded {
return Err(CoreError::sys(libc::EOVERFLOW, "spawn output limit"));
}
Ok(Output {
pid,
status: self.status.take(),
stdout,
stderr,
timed_out: self.timed_out,
stdout_early_exited,
})
}
}
impl Drop for ManagedProcess {
fn drop(&mut self) {
let Some(running) = self.running.take() else {
return;
};
let process = &running.process;
let pid = process.pid();
let pgid = effective_pgid(pid, self.pgroup);
let target_is_group = self.pgroup.isolated || self.pgroup.leader.is_some();
let _ = signal_process(process, target_is_group, pgid, libc::SIGKILL);
let _ = process.wait_blocking();
}
}
fn effective_pgid(pid: pid_t, pgroup: ProcessGroup) -> pid_t {
match pgroup.leader {
Some(0) | None => pid,
Some(leader) => leader,
}
}
fn signal_process(
process: &Process,
target_is_group: bool,
pgid: pid_t,
signal: i32,
) -> Result<(), CoreError> {
if target_is_group {
process.kill_group(pgid, signal)
} else {
process.kill(signal)
}
}
pub fn spawn_start(opts: SpawnOptions) -> Result<RunningProcess, CoreError> {
if !opts.wait && (opts.stdin.is_some() || opts.capture_stdout || opts.capture_stderr) {
return Err(CoreError::sys(
libc::EINVAL,
"background I/O capture not supported (wait must be true)",
));
}
validate_backend(&opts)?;
let (pid, drain) = match opts.backend {
SpawnBackend::PosixSpawn => spawn_posix_internal(opts)?,
SpawnBackend::Fork => spawn_fork_internal(opts)?,
};
Ok(RunningProcess {
process: Process::new(pid),
drain,
})
}
pub fn spawn_managed(opts: SpawnOptions) -> Result<ManagedProcess, CoreError> {
if !opts.wait {
return Err(CoreError::sys(
libc::EINVAL,
"managed process requires wait=true",
));
}
let timeout_at = opts
.timeout_ms
.map(|ms| Instant::now() + Duration::from_millis(ms as u64));
let kill_grace = Duration::from_millis(opts.kill_grace_ms as u64);
let cancel = opts.cancel;
let pgroup = opts.pgroup;
let running = spawn_start(opts)?;
Ok(ManagedProcess {
running: Some(running),
timeout_at,
kill_grace,
cancel,
pgroup,
cancel_at: None,
kill_state: KillState::None,
status: None,
timed_out: false,
})
}
pub fn spawn(opts: SpawnOptions) -> Result<Output, CoreError> {
let wait = opts.wait;
let timeout_ms = opts.timeout_ms;
let kill_grace_ms = opts.kill_grace_ms;
let cancel = opts.cancel;
let pgroup = opts.pgroup;
let mut reactor = Reactor::new()?;
let running = spawn_start(opts)?;
let pid = running.process.pid();
let mut drain = running.drain;
drain.register_with_reactor(&mut reactor)?;
if !wait {
let (stdout, stderr) = drain.into_parts();
return Ok(Output {
pid,
status: None,
stdout,
stderr,
timed_out: false,
stdout_early_exited: false,
});
}
wait_loop(
pid,
drain,
reactor,
timeout_ms,
kill_grace_ms,
cancel,
pgroup,
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum KillState {
None,
TermSent,
KillSent,
}
fn wait_loop(
pid: pid_t,
mut drain: crate::io::DrainState<fn(&[u8]) -> bool>,
mut reactor: Reactor,
timeout_ms: Option<u32>,
kill_grace_ms: u32,
cancel: CancelPolicy,
pgroup: ProcessGroup,
) -> Result<Output, CoreError> {
let process = Process::new(pid);
let pgid = effective_pgid(pid, pgroup);
let mut status_raw = process.wait_step()?;
let mut state = KillState::None;
let mut timed_out = false;
let start_time = std::time::Instant::now();
let deadline = timeout_ms.map(|t| std::time::Duration::from_millis(t as u64));
loop {
let mut poll_timeout = -1;
if let Some(dl) = deadline {
let elapsed = start_time.elapsed();
if elapsed >= dl {
timed_out = true;
let elapsed_over = (elapsed - dl).as_millis();
let target_is_group = pgroup.isolated || pgroup.leader.is_some();
match state {
KillState::None => {
if cancel == CancelPolicy::Graceful {
let r = if target_is_group {
process.kill_group(pgid, libc::SIGTERM)
} else {
process.kill(libc::SIGTERM)
};
if r.is_err() {
state = KillState::KillSent; } else {
state = KillState::TermSent;
}
} else if cancel == CancelPolicy::Kill {
let _ = if target_is_group {
process.kill_group(pgid, libc::SIGKILL)
} else {
process.kill(libc::SIGKILL)
};
state = KillState::KillSent;
} else {
}
}
KillState::TermSent if elapsed_over > kill_grace_ms as u128 => {
let _ = if target_is_group {
process.kill_group(pgid, libc::SIGKILL)
} else {
process.kill(libc::SIGKILL)
};
state = KillState::KillSent;
}
_ => {}
}
poll_timeout = 100; } else {
let remaining = dl - elapsed;
poll_timeout = remaining.as_millis().min(i32::MAX as u128) as i32;
}
}
if status_raw.is_none()
&& let Some(s) = process.wait_step()?
{
status_raw = Some(s);
}
if drain.is_done() {
let s = if status_raw.is_some() {
status_raw.take()
} else if deadline.is_none() {
Some(process.wait_blocking()?)
} else {
None
};
if let Some(s) = s {
for slot in drain.take_all_slots() {
reactor.del(&slot.fd)?;
}
let (stdout, stderr, output_limit_exceeded, stdout_early_exited) =
drain.into_parts_with_state();
if output_limit_exceeded {
return Err(CoreError::sys(libc::EOVERFLOW, "spawn output limit"));
}
return Ok(Output {
pid,
status: Some(s),
stdout,
stderr,
timed_out,
stdout_early_exited,
});
}
}
if timed_out && status_raw.is_some() {
for slot in drain.take_all_slots() {
let _ = reactor.del(&slot.fd);
}
let (stdout, stderr, _output_limit_exceeded, stdout_early_exited) =
drain.into_parts_with_state();
return Ok(Output {
pid,
status: status_raw,
stdout,
stderr,
timed_out: true,
stdout_early_exited,
});
}
let timeout = poll_timeout;
let mut events = Vec::new();
let nevents = reactor.wait(&mut events, 64, timeout)?;
for ev in events.iter().take(nevents) {
if drain.stdout_matches(ev.token) {
if ev.readable || ev.hangup {
drain.handle_stdout_ready(&mut reactor)?;
} else if ev.error {
drain.drop_stdout(&mut reactor)?;
}
} else if drain.stderr_matches(ev.token) {
if ev.readable || ev.hangup {
drain.handle_stderr_ready(&mut reactor)?;
} else if ev.error {
drain.drop_stderr(&mut reactor)?;
}
} else if drain.stdin_matches(ev.token) {
if ev.writable {
drain.handle_stdin_writable(&mut reactor)?;
} else if ev.error || ev.hangup {
drain.drop_stdin(&mut reactor)?;
}
}
}
}
}