use std::collections::VecDeque;
use std::io::Read;
use std::process::{Child, ChildStdin, Command, Stdio};
use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread;
use std::time::{Duration, Instant};
use crate::observer::{ObserverEmitter, ProcessWatchEmitter};
pub use running_process_platform_internal::foreground;
pub(crate) use running_process_platform_internal::platform;
#[cfg(feature = "async-process")]
mod async_process;
#[cfg(feature = "async-process")]
mod blocking_island;
#[cfg(feature = "async-process")]
pub use blocking_island::dispatch_blocking as blocking_island_dispatch;
pub mod console_detect;
pub mod containment;
mod descendant_monitor;
pub mod env_vars;
pub mod environment;
mod helpers;
#[cfg(feature = "async-process")]
mod process_runtime;
#[cfg(feature = "window-icon")]
pub mod window_icon;
#[cfg(feature = "daemon-registration")]
pub mod daemon_registration;
#[cfg(feature = "daemon-registration")]
pub mod daemon_registration_compat;
#[cfg(feature = "daemon-registration-v2")]
pub mod daemon_registration_v2;
#[cfg(feature = "daemon-registration-v2")]
pub mod daemon_registration_v2_compat;
#[cfg(feature = "frame-v1-codec")]
pub mod daemon_frame_v1;
#[cfg(any(feature = "daemon-registration", feature = "daemon-registration-v2"))]
pub(crate) mod daemon_registration_common;
#[cfg(feature = "frame-v1-codec")]
pub mod frame_v1;
#[cfg(any(feature = "backend-identity", feature = "daemon-registration"))]
#[path = "broker/host_identity.rs"]
pub(crate) mod daemon_host_identity;
pub mod observer;
#[cfg(feature = "originator-scan")]
pub mod originator;
pub mod output_log;
#[cfg(feature = "client")]
pub mod proto {
pub use running_process_protocol::daemon;
}
#[cfg(feature = "client")]
pub mod client;
#[cfg(feature = "client")]
pub mod broker;
#[cfg(all(feature = "backend-identity", not(feature = "client")))]
#[allow(dead_code, unused_imports)]
mod broker;
#[cfg(feature = "backend-identity")]
pub mod backend_identity;
#[cfg(feature = "client")]
pub mod content_hash;
#[cfg(all(feature = "backend-identity", not(feature = "client")))]
mod content_hash;
#[cfg(feature = "probe")]
pub mod probe;
#[cfg(feature = "client")]
pub mod maintenance;
#[cfg(feature = "client")]
pub mod cleanup;
#[cfg(feature = "client")]
pub mod boot_autostart;
#[cfg(feature = "client")]
pub mod runpm_config;
#[cfg(feature = "test-support")]
pub mod test_support;
#[cfg(all(feature = "telemetry", not(feature = "daemon")))]
#[path = "daemon/telemetry.rs"]
pub mod telemetry;
#[cfg(all(feature = "telemetry", feature = "daemon"))]
pub use daemon::telemetry;
#[cfg(all(feature = "telemetry", feature = "daemon"))]
const _: fn(crate::telemetry::TeeHandle) -> daemon::telemetry::TeeHandle = |handle| handle;
#[cfg(feature = "daemon")]
pub mod daemon;
#[cfg(feature = "independent-spawn")]
pub mod independent_spawn;
pub mod process_tree;
#[cfg(feature = "pty")]
pub mod pty;
mod public_symbols;
mod rust_debug;
pub mod spawn;
mod spawn_contract;
pub use spawn_contract::{IndependentBackend, SpawnLifetime, SpawnMode, SpawnOptions};
#[cfg(feature = "independent-spawn")]
mod spawn_dispatch;
#[cfg(feature = "independent-spawn")]
pub use spawn_dispatch::{spawn_with_options, SpawnExit, SpawnHandle};
pub mod systemd_killmode;
#[cfg(feature = "terminal-graphics")]
pub mod terminal_graphics;
mod types;
#[cfg(unix)]
mod unix;
#[cfg(windows)]
mod windows;
#[cfg(feature = "async-process")]
pub use async_process::{
AsyncCapturedOutput, AsyncProcess, AsyncProcessBuilder, AsyncProcessSession,
AsyncProcessSessionChunk, AsyncProcessSessionControl, AsyncProcessSessionEvent,
AsyncProcessSessionOptions, AsyncProcessSessionOutput, AsyncStdio, ProcessTreeKill,
};
pub use console_detect::{monitor_console_windows, ConsoleWindowInfo};
pub use containment::{ContainedProcessGroup, ORIGINATOR_ENV_VAR};
#[cfg(feature = "client")]
pub use content_hash::blake3_file;
pub use observer::{
CapabilitySupport, CaptureSource, CategoryCapability, DumpResult, EventCategory,
ObservationGrade, ObservationPolicy, ObserverCapabilities, ObserverConfig, ObserverEvent,
ObserverEventKind, ObserverSubscriber, ProcessEvent, ProcessEventKind, ProcessIdentity,
ProcessObservation, ProcessObservationCapabilities, ProcessObservationError, ProcessWatch,
ProcessWatchConfigurationError, ProcessWatchCursor, ProcessWatchGap, ProcessWatchLoss,
ProcessWatchMatch, ProcessWatchRead, ProcessWatchSubscriber, StackCapture, StackDump,
};
#[cfg(feature = "originator-scan")]
pub use originator::{
find_declared_daemon_pids, find_processes_by_originator, OriginatorProcessInfo,
};
pub use output_log::{
CursorRead, OutputCursor, OutputLog, OutputRecord, SharedOutputCursor, SharedOutputLog,
};
#[doc(hidden)]
pub use running_process_platform_internal::platform::executable as platform_executable;
#[cfg(target_os = "linux")]
pub use running_process_platform_internal::platform::process::current_executable_build_id;
pub use running_process_platform_internal::platform::process::{
ProcessInspectError, ProcessInspectErrorKind,
};
pub use running_process_platform_internal::process_executable_path;
pub use running_process_platform_internal::process_same_executable_path;
pub use running_process_platform_internal::ProcessLiveness;
pub use rust_debug::{render_rust_debug_traces, RustDebugScopeGuard};
pub use spawn::{
spawn, spawn_daemon, spawn_daemon_breaking_away_from_job,
spawn_daemon_breaking_away_with_env_policy, spawn_daemon_with_clear_env,
spawn_daemon_with_env_policy, spawn_daemon_with_environment,
spawn_daemon_with_explicit_environment, spawn_daemon_with_stdio,
spawn_daemon_with_stdio_and_env_policy, spawn_with_env_policy, spawn_with_environment,
spawn_with_explicit_environment, DaemonChild, DaemonStdio, DaemonStdioSource,
EnvironmentPolicy, SpawnStdio, SpawnedChild, SpawnedChildControl, StdioSource, SyncEnvironment,
DAEMON_MARKER_ENV_VAR,
};
#[cfg(feature = "client-async")]
pub use spawn::{spawn_tokio, TokioSpawnOptions};
#[cfg(feature = "terminal-graphics")]
pub use terminal_graphics::{
current_terminal_capabilities, current_terminal_capabilities_with_timeout,
detect_terminal_capabilities, CapabilityStatus, EvidenceStrength, GraphicsCapability,
GraphicsProtocol, TerminalCapabilities, TerminalCapabilityInput, TerminalGraphicsCapabilities,
TerminalProbeEvidence,
};
pub use types::{
CommandSpec, ProcessConfig, ProcessError, ReadStatus, RunOutput, StderrMode, StdinMode,
StreamEvent, StreamKind,
};
#[cfg(feature = "window-icon")]
pub use window_icon::{
host_icon_support, icon_support, set_host_icon, set_icon, IconError, IconScope, IconSource,
IconSupport, StockIcon,
};
#[cfg(unix)]
pub(crate) use helpers::{child_try_wait_error_is_retryable, poll_mutex_until};
pub(crate) use helpers::{exit_code, feed_chunk, kill_drain_deadline, log_spawned_child_pid};
pub use running_process_platform_internal::exit_code as native_exit_code;
pub use running_process_platform_internal::ProcessPriority;
#[cfg(feature = "async-process")]
pub use running_process_platform_internal::SpawnAdmission;
#[cfg(unix)]
pub use unix::{unix_set_priority, unix_signal_process, unix_signal_process_group, UnixSignal};
#[cfg(windows)]
pub(crate) use windows::{assign_child_to_windows_kill_on_close_job_impl, WindowsJobHandle};
#[macro_export]
macro_rules! rp_rust_debug_scope {
($label:expr) => {
let _running_process_rust_debug_scope =
$crate::RustDebugScopeGuard::enter($label, file!(), line!());
};
}
#[derive(Default)]
struct QueueState {
stdout_queue: VecDeque<Vec<u8>>,
stderr_queue: VecDeque<Vec<u8>>,
combined_queue: VecDeque<StreamEvent>,
stdout_history: VecDeque<Vec<u8>>,
stderr_history: VecDeque<Vec<u8>>,
combined_history: VecDeque<StreamEvent>,
stdout_raw: VecDeque<Vec<u8>>,
stderr_raw: VecDeque<Vec<u8>>,
stdout_raw_bytes: usize,
stderr_raw_bytes: usize,
stdout_history_bytes: usize,
stderr_history_bytes: usize,
combined_history_bytes: usize,
stdout_closed: bool,
stderr_closed: bool,
}
const RETURNCODE_NOT_SET: i64 = i64::MIN;
struct SharedState {
queues: Mutex<QueueState>,
condvar: Condvar,
capture_limit: Option<usize>,
capture_overflowed: AtomicBool,
active_capture_readers: std::sync::atomic::AtomicUsize,
returncode: AtomicI64,
observer: Option<ObserverEmitter>,
observer_exit_emitted: AtomicBool,
}
struct ChildState {
child: ChildHandle,
#[cfg(windows)]
_job: WindowsJobHandle,
}
enum ChildHandle {
Standard(Child),
ExactTrace(running_process_platform_internal::platform::process::TracedChild),
}
impl ChildHandle {
fn id(&self) -> u32 {
match self {
Self::Standard(child) => child.id(),
Self::ExactTrace(child) => child.id(),
}
}
fn try_wait_code(&mut self) -> std::io::Result<Option<i32>> {
match self {
Self::Standard(child) => child.try_wait().map(|status| status.map(exit_code)),
Self::ExactTrace(child) => child.try_wait_code(),
}
}
fn kill(&mut self) -> std::io::Result<()> {
match self {
Self::Standard(child) => child.kill(),
Self::ExactTrace(child) => child.kill(),
}
}
fn take_stdin(&mut self) -> Option<ChildStdin> {
match self {
Self::Standard(child) => child.stdin.take(),
Self::ExactTrace(child) => child.take_stdin(),
}
}
fn take_stdout(&mut self) -> Option<std::process::ChildStdout> {
match self {
Self::Standard(child) => child.stdout.take(),
Self::ExactTrace(child) => child.take_stdout(),
}
}
fn take_stderr(&mut self) -> Option<std::process::ChildStderr> {
match self {
Self::Standard(child) => child.stderr.take(),
Self::ExactTrace(child) => child.take_stderr(),
}
}
#[cfg(windows)]
fn wait_code(&mut self) -> std::io::Result<i32> {
match self {
Self::Standard(child) => child.wait().map(exit_code),
Self::ExactTrace(child) => child.wait_code(),
}
}
}
#[cfg(test)]
#[derive(Debug, Eq, PartialEq)]
enum CapturePollAction {
Wait,
Read,
Cancel,
}
#[cfg(test)]
fn capture_poll_action(capture_revents: i16, wake_revents: i16) -> CapturePollAction {
if wake_revents != 0 {
CapturePollAction::Cancel
} else if capture_revents != 0 {
CapturePollAction::Read
} else {
CapturePollAction::Wait
}
}
fn cleanup_child_after_start_error(child: ChildHandle) {
match child {
ChildHandle::Standard(mut child) => {
let _ = child.kill();
thread::spawn(move || {
let _ = child.wait();
});
}
ChildHandle::ExactTrace(mut child) => {
let _ = child.kill();
}
}
}
impl SharedState {
#[cfg(test)]
fn new(capture: bool) -> Self {
Self::with_observer_and_limit(capture, None, None)
}
fn with_observer_and_limit(
capture: bool,
observer: Option<ObserverEmitter>,
capture_limit: Option<usize>,
) -> Self {
let queues = QueueState {
stdout_closed: !capture,
stderr_closed: !capture,
..QueueState::default()
};
Self {
queues: Mutex::new(queues),
condvar: Condvar::new(),
capture_limit,
capture_overflowed: AtomicBool::new(false),
active_capture_readers: std::sync::atomic::AtomicUsize::new(0),
returncode: AtomicI64::new(RETURNCODE_NOT_SET),
observer,
observer_exit_emitted: AtomicBool::new(false),
}
}
fn emit_exited(&self, pid: u32, exit_code: i32) {
let Some(emitter) = self.observer.as_ref() else {
return;
};
if self
.observer_exit_emitted
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
emitter.emit_exited(pid, exit_code);
}
}
}
pub struct NativeProcess {
config: ProcessConfig,
command_override: Mutex<Option<Command>>,
child: Arc<Mutex<Option<ChildState>>>,
stdin: Mutex<Option<ChildStdin>>,
shared: Arc<SharedState>,
process_watch: Option<Arc<ProcessWatchEmitter>>,
kill_when_owner_dies: bool,
#[cfg(test)]
stdin_write_active: AtomicBool,
capture_cancellation:
Arc<running_process_platform_internal::platform::process::CaptureCancellation>,
}
impl NativeProcess {
pub fn new(config: ProcessConfig) -> Self {
Self::new_with_options(config, None, None, None, None)
}
pub fn with_observer(
config: ProcessConfig,
observer: crate::observer::ObserverConfig,
) -> (Self, ObserverSubscriber) {
let (emitter, subscriber) = ObserverEmitter::new(observer);
let process = Self::new_with_options(config, Some(emitter), None, None, None);
(process, subscriber)
}
pub fn with_observer_and_command(
command: Command,
config: ProcessConfig,
observer: crate::observer::ObserverConfig,
) -> (Self, ObserverSubscriber) {
let (emitter, subscriber) = ObserverEmitter::new(observer);
let process = Self::new_with_options(config, Some(emitter), None, Some(command), None);
(process, subscriber)
}
pub fn with_process_watches(
config: ProcessConfig,
watches: Vec<ProcessWatch>,
policy: ObservationPolicy,
) -> Result<(Self, ProcessWatchSubscriber), ProcessObservationError> {
let (emitter, subscriber) = ProcessWatchEmitter::new(watches, policy)?;
let process = Self::new_with_options(config, None, None, None, Some(emitter));
Ok((process, subscriber))
}
pub fn process_observation_capabilities() -> ProcessObservationCapabilities {
ProcessObservationCapabilities::current()
}
fn new_with_capture_limit(config: ProcessConfig, capture_limit: usize) -> Self {
Self::new_with_options(config, None, Some(capture_limit), None, None)
}
fn new_with_command_capture_limit(
command: Command,
config: ProcessConfig,
capture_limit: usize,
kill_when_owner_dies: bool,
) -> Self {
let mut process =
Self::new_with_options(config, None, Some(capture_limit), Some(command), None);
process.kill_when_owner_dies = kill_when_owner_dies;
process
}
fn new_with_options(
config: ProcessConfig,
observer: Option<ObserverEmitter>,
capture_limit: Option<usize>,
command_override: Option<Command>,
process_watch: Option<Arc<ProcessWatchEmitter>>,
) -> Self {
let shared = SharedState::with_observer_and_limit(config.capture, observer, capture_limit);
Self {
shared: Arc::new(shared),
process_watch,
command_override: Mutex::new(command_override),
child: Arc::new(Mutex::new(None)),
stdin: Mutex::new(None),
kill_when_owner_dies: false,
#[cfg(test)]
stdin_write_active: AtomicBool::new(false),
config,
capture_cancellation: Arc::new(Default::default()),
}
}
#[inline(never)]
pub fn start(&self) -> Result<(), ProcessError> {
public_symbols::rp_native_process_start_public(self)
}
fn start_impl(&self) -> Result<(), ProcessError> {
crate::rp_rust_debug_scope!("running_process::NativeProcess::start");
let mut guard = self.child.lock().expect("child mutex poisoned");
if guard.is_some() {
return Err(ProcessError::AlreadyStarted);
}
let mut command = self.build_command();
let exact_trace = self
.process_watch
.as_ref()
.is_some_and(|watch| watch.uses_exact_trace());
match self.config.stdin_mode {
StdinMode::Inherit => {}
StdinMode::Piped => {
command.stdin(Stdio::piped());
}
StdinMode::Null => {
command.stdin(Stdio::null());
}
}
if self.config.capture {
command.stdout(Stdio::piped());
command.stderr(Stdio::piped());
}
let mut child = if exact_trace {
let event_watch = Arc::clone(self.process_watch.as_ref().expect("exact watch checked"));
let completion_watch = Arc::clone(&event_watch);
match running_process_platform_internal::platform::process::start_exact_trace(
command,
Box::new(move |event| event_watch.emit_exact(event)),
Box::new(move || completion_watch.close()),
) {
Ok(child) => ChildHandle::ExactTrace(child),
Err(error) => {
if let Some(watch) = self.process_watch.as_ref() {
watch.close();
}
return Err(ProcessError::Spawn(error));
}
}
} else {
ChildHandle::Standard(command.spawn().map_err(ProcessError::Spawn)?)
};
log_spawned_child_pid(child.id()).map_err(ProcessError::Spawn)?;
if let Some(emitter) = self.shared.observer.as_ref() {
emitter.emit_started(child.id());
}
#[cfg(windows)]
let job = {
let descendant_sink = self
.shared
.observer
.as_ref()
.and_then(|e| e.descendant_sink());
let job_result = match &child {
ChildHandle::Standard(standard_child) => {
let direct_pid = standard_child.id();
public_symbols::rp_assign_child_to_windows_kill_on_close_job_with_observer_public(
standard_child,
descendant_sink,
self.process_watch.clone(),
direct_pid,
self.config.address_space_limit_bytes,
)
}
ChildHandle::ExactTrace(_) => unreachable!("Windows exact tracing is unavailable"),
};
match job_result {
Ok(job) => job,
Err(error) => {
if let Some(watch) = self.process_watch.as_ref() {
watch.close();
}
cleanup_child_after_start_error(child);
return Err(ProcessError::Spawn(error));
}
}
};
if !exact_trace {
descendant_monitor::start(
child.id(),
self.shared.observer.as_ref(),
self.process_watch.as_ref(),
);
}
if self.config.capture {
let stdout = child.take_stdout().expect("stdout pipe missing");
let stderr = child.take_stderr().expect("stderr pipe missing");
let stdout =
match running_process_platform_internal::platform::process::prepare_capture_reader(
stdout,
&self.capture_cancellation,
running_process_platform_internal::platform::process::CaptureStream::Stdout,
) {
Ok(stdout) => stdout,
Err(error) => {
cleanup_child_after_start_error(child);
return Err(ProcessError::Spawn(error));
}
};
let stderr =
match running_process_platform_internal::platform::process::prepare_capture_reader(
stderr,
&self.capture_cancellation,
running_process_platform_internal::platform::process::CaptureStream::Stderr,
) {
Ok(stderr) => stderr,
Err(error) => {
running_process_platform_internal::platform::process::capture_reader_done(
&self.capture_cancellation,
running_process_platform_internal::platform::process::CaptureStream::Stdout,
);
cleanup_child_after_start_error(child);
return Err(ProcessError::Spawn(error));
}
};
self.spawn_reader(
stdout,
StreamKind::Stdout,
StreamKind::Stdout,
self.pipe_done_callback(StreamKind::Stdout),
);
self.spawn_reader(
stderr,
StreamKind::Stderr,
match self.config.stderr_mode {
StderrMode::Stdout => StreamKind::Stdout,
StderrMode::Pipe => StreamKind::Stderr,
},
self.pipe_done_callback(StreamKind::Stderr),
);
}
*self.stdin.lock().expect("stdin mutex poisoned") = child.take_stdin();
*guard = Some(ChildState {
child,
#[cfg(windows)]
_job: job,
});
drop(guard);
self.spawn_exit_waiter();
Ok(())
}
fn spawn_exit_waiter(&self) {
let child = Arc::clone(&self.child);
let shared = Arc::clone(&self.shared);
let capture = self.config.capture;
let capture_cancellation = Arc::clone(&self.capture_cancellation);
thread::spawn(move || {
loop {
if shared.returncode.load(Ordering::Acquire) != RETURNCODE_NOT_SET {
return;
}
let exited = {
let mut guard = child.lock().expect("child mutex poisoned");
if let Some(child_state) = guard.as_mut() {
let pid = child_state.child.id();
match child_state.child.try_wait_code() {
Ok(Some(code)) => {
shared.returncode.store(code as i64, Ordering::Release);
shared.emit_exited(pid, code);
shared.condvar.notify_all();
true
}
Ok(None) => false,
Err(_error) => {
#[cfg(unix)]
if child_try_wait_error_is_retryable(&_error) {
false
} else {
return;
}
#[cfg(windows)]
return;
}
}
} else {
return;
}
};
if exited {
if capture {
let drained = finalize_capture_completion(&shared, kill_drain_deadline());
if !drained {
running_process_platform_internal::platform::process::cancel_capture_reader(
&capture_cancellation,
);
}
}
return;
}
thread::sleep(Duration::from_millis(10));
}
});
}
pub fn write_stdin(&self, data: &[u8]) -> Result<(), ProcessError> {
if self.child.lock().expect("child mutex poisoned").is_none() {
return Err(ProcessError::NotRunning);
}
let mut guard = self.stdin.lock().expect("stdin mutex poisoned");
let stdin = guard.as_mut().ok_or(ProcessError::StdinUnavailable)?;
use std::io::Write;
#[cfg(test)]
self.stdin_write_active.store(true, Ordering::Release);
let write_result = stdin.write_all(data);
#[cfg(test)]
self.stdin_write_active.store(false, Ordering::Release);
write_result.map_err(ProcessError::Io)?;
stdin.flush().map_err(ProcessError::Io)?;
drop(guard.take());
Ok(())
}
pub fn write_stdin_streaming(&self, data: &[u8]) -> Result<(), ProcessError> {
if self.child.lock().expect("child mutex poisoned").is_none() {
return Err(ProcessError::NotRunning);
}
let mut guard = self.stdin.lock().expect("stdin mutex poisoned");
let stdin = guard.as_mut().ok_or(ProcessError::StdinUnavailable)?;
use std::io::Write;
#[cfg(test)]
self.stdin_write_active.store(true, Ordering::Release);
let write_result = stdin.write_all(data);
#[cfg(test)]
self.stdin_write_active.store(false, Ordering::Release);
write_result.map_err(ProcessError::Io)?;
stdin.flush().map_err(ProcessError::Io)?;
Ok(())
}
pub fn close_stdin(&self) -> Result<(), ProcessError> {
if self.child.lock().expect("child mutex poisoned").is_none() {
return Err(ProcessError::NotRunning);
}
drop(self.stdin.lock().expect("stdin mutex poisoned").take());
Ok(())
}
pub fn poll(&self) -> Result<Option<i32>, ProcessError> {
if let Some(code) = self.returncode() {
return Ok(Some(code));
}
let mut guard = self.child.lock().expect("child mutex poisoned");
let Some(child_state) = guard.as_mut() else {
return Ok(self.returncode());
};
let pid = child_state.child.id();
let child = &mut child_state.child;
let status = child.try_wait_code().map_err(ProcessError::Io)?;
if let Some(code) = status {
self.set_returncode(code);
self.shared.emit_exited(pid, code);
return Ok(Some(code));
}
Ok(None)
}
#[inline(never)]
pub fn wait(&self, timeout: Option<Duration>) -> Result<i32, ProcessError> {
public_symbols::rp_native_process_wait_public(self, timeout)
}
fn wait_impl(&self, timeout: Option<Duration>) -> Result<i32, ProcessError> {
crate::rp_rust_debug_scope!("running_process::NativeProcess::wait");
if self.child.lock().expect("child mutex poisoned").is_none() {
return self.returncode().ok_or(ProcessError::NotRunning);
}
if let Some(code) = self.returncode() {
self.finish_capture_drain();
return Ok(code);
}
let start = Instant::now();
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
loop {
let rc = self.shared.returncode.load(Ordering::Acquire);
if rc != RETURNCODE_NOT_SET {
drop(guard);
let code = rc as i32;
self.finish_capture_drain();
return Ok(code);
}
if let Some(limit) = timeout {
let elapsed = start.elapsed();
if elapsed >= limit {
return Err(ProcessError::Timeout);
}
let remaining = limit - elapsed;
let wait_time = remaining.min(Duration::from_millis(50));
guard = self
.shared
.condvar
.wait_timeout(guard, wait_time)
.expect("queue mutex poisoned")
.0;
} else {
guard = self
.shared
.condvar
.wait_timeout(guard, Duration::from_millis(50))
.expect("queue mutex poisoned")
.0;
}
}
}
#[inline(never)]
pub fn kill(&self) -> Result<(), ProcessError> {
public_symbols::rp_native_process_kill_public(self)
}
fn kill_impl(&self) -> Result<(), ProcessError> {
crate::rp_rust_debug_scope!("running_process::NativeProcess::kill");
#[cfg(windows)]
{
let mut guard = self.child.lock().expect("child mutex poisoned");
let child = &mut guard.as_mut().ok_or(ProcessError::NotRunning)?.child;
let pid = child.id();
child.kill().map_err(ProcessError::Io)?;
let code = child.wait_code().map_err(ProcessError::Io)?;
self.set_returncode(code);
self.shared.emit_exited(pid, code);
}
#[cfg(unix)]
{
let deadline = kill_drain_deadline();
let (pid, already_reaped) = {
let mut state = self.child.lock().expect("child mutex poisoned");
let child = &mut state.as_mut().ok_or(ProcessError::NotRunning)?.child;
let pid = child.id();
if let Some(code) = child.try_wait_code().map_err(ProcessError::Io)? {
(pid, Some(code))
} else {
let group_signaled = self.config.create_process_group
&& unix_signal_process_group(pid as i32, UnixSignal::Kill).is_ok();
if !group_signaled {
child.kill().map_err(ProcessError::Io)?;
}
(pid, None)
}
};
self.cancel_capture_io();
let reaped = already_reaped.or_else(|| {
let reap_result =
poll_mutex_until(&self.child, deadline, Duration::from_millis(10), |state| {
match state.as_mut() {
Some(child) => child.child.try_wait_code(),
None => Ok(None),
}
});
match reap_result {
Ok(Some(code)) => Some(code),
_ => None,
}
});
if let Some(code) = reaped {
self.set_returncode(code);
self.shared.emit_exited(pid, code);
}
public_symbols::rp_native_process_wait_for_capture_completion_with_deadline_public(
self, deadline,
);
Ok(())
}
#[cfg(windows)]
{
self.cancel_capture_io();
public_symbols::rp_native_process_wait_for_capture_completion_with_deadline_public(
self,
kill_drain_deadline(),
);
Ok(())
}
}
pub fn terminate(&self) -> Result<(), ProcessError> {
self.kill()
}
pub fn terminate_group_soft(&self) -> Result<(), ProcessError> {
if !self.config.create_process_group {
return Ok(());
}
let pid = self.pid().ok_or(ProcessError::NotRunning)?;
running_process_platform_internal::platform::process::soft_terminate_process_group(pid)
.map_err(ProcessError::Io)
}
#[inline(never)]
pub fn close(&self) -> Result<(), ProcessError> {
public_symbols::rp_native_process_close_public(self)
}
fn close_impl(&self) -> Result<(), ProcessError> {
crate::rp_rust_debug_scope!("running_process::NativeProcess::close");
if self.child.lock().expect("child mutex poisoned").is_none() {
return Ok(());
}
if self.poll()?.is_none() {
self.kill()?;
} else {
self.finish_capture_drain();
}
if let Some(watch) = self.process_watch.as_ref() {
watch.close();
}
Ok(())
}
pub fn pid(&self) -> Option<u32> {
self.child
.lock()
.expect("child mutex poisoned")
.as_ref()
.map(|state| state.child.id())
}
pub fn returncode(&self) -> Option<i32> {
let v = self.shared.returncode.load(Ordering::Acquire);
if v == RETURNCODE_NOT_SET {
None
} else {
Some(v as i32)
}
}
pub fn has_pending_stream(&self, stream: StreamKind) -> bool {
if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
return false;
}
let guard = self.shared.queues.lock().expect("queue mutex poisoned");
match stream {
StreamKind::Stdout => !guard.stdout_queue.is_empty(),
StreamKind::Stderr => !guard.stderr_queue.is_empty(),
}
}
pub fn has_pending_combined(&self) -> bool {
let guard = self.shared.queues.lock().expect("queue mutex poisoned");
!guard.combined_queue.is_empty()
}
pub fn drain_stream(&self, stream: StreamKind) -> Vec<Vec<u8>> {
if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
return Vec::new();
}
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
let queue = match stream {
StreamKind::Stdout => &mut guard.stdout_queue,
StreamKind::Stderr => &mut guard.stderr_queue,
};
queue.drain(..).collect()
}
pub fn drain_stream_raw(&self, stream: StreamKind) -> Vec<u8> {
if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
return Vec::new();
}
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
match stream {
StreamKind::Stdout => {
let mut output = Vec::with_capacity(guard.stdout_raw_bytes);
for chunk in guard.stdout_raw.drain(..) {
output.extend_from_slice(&chunk);
}
guard.stdout_raw_bytes = 0;
output
}
StreamKind::Stderr => {
let mut output = Vec::with_capacity(guard.stderr_raw_bytes);
for chunk in guard.stderr_raw.drain(..) {
output.extend_from_slice(&chunk);
}
guard.stderr_raw_bytes = 0;
output
}
}
}
pub fn drain_combined(&self) -> Vec<StreamEvent> {
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
guard.combined_queue.drain(..).collect()
}
pub fn read_stream(
&self,
stream: StreamKind,
timeout: Option<Duration>,
) -> ReadStatus<Vec<u8>> {
let deadline = timeout.map(|limit| Instant::now() + limit);
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
loop {
if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
return ReadStatus::Eof;
}
let queue = match stream {
StreamKind::Stdout => &mut guard.stdout_queue,
StreamKind::Stderr => &mut guard.stderr_queue,
};
if let Some(line) = queue.pop_front() {
return ReadStatus::Line(line);
}
let closed = match stream {
StreamKind::Stdout => {
if self.config.stderr_mode == StderrMode::Stdout {
guard.stdout_closed && guard.stderr_closed
} else {
guard.stdout_closed
}
}
StreamKind::Stderr => guard.stderr_closed,
};
if closed {
return ReadStatus::Eof;
}
match deadline {
Some(deadline) => {
let now = Instant::now();
if now >= deadline {
return ReadStatus::Timeout;
}
let wait = deadline.saturating_duration_since(now);
let result = self
.shared
.condvar
.wait_timeout(guard, wait)
.expect("queue mutex poisoned");
guard = result.0;
if result.1.timed_out() {
return ReadStatus::Timeout;
}
}
None => {
guard = self
.shared
.condvar
.wait(guard)
.expect("queue mutex poisoned");
}
}
}
}
#[inline(never)]
pub fn read_combined(&self, timeout: Option<Duration>) -> ReadStatus<StreamEvent> {
public_symbols::rp_native_process_read_combined_public(self, timeout)
}
fn read_combined_impl(&self, timeout: Option<Duration>) -> ReadStatus<StreamEvent> {
crate::rp_rust_debug_scope!("running_process::NativeProcess::read_combined");
let deadline = timeout.map(|limit| Instant::now() + limit);
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
loop {
if let Some(event) = guard.combined_queue.pop_front() {
return ReadStatus::Line(event);
}
if guard.stdout_closed && guard.stderr_closed {
return ReadStatus::Eof;
}
match deadline {
Some(deadline) => {
let now = Instant::now();
if now >= deadline {
return ReadStatus::Timeout;
}
let wait = deadline.saturating_duration_since(now);
let result = self
.shared
.condvar
.wait_timeout(guard, wait)
.expect("queue mutex poisoned");
guard = result.0;
if result.1.timed_out() {
return ReadStatus::Timeout;
}
}
None => {
guard = self
.shared
.condvar
.wait(guard)
.expect("queue mutex poisoned");
}
}
}
}
pub fn captured_stdout(&self) -> Vec<Vec<u8>> {
self.shared
.queues
.lock()
.expect("queue mutex poisoned")
.stdout_history
.clone()
.into_iter()
.collect()
}
fn captured_stdout_raw(&self) -> Vec<u8> {
let guard = self.shared.queues.lock().expect("queue mutex poisoned");
guard.stdout_raw.iter().flatten().copied().collect()
}
pub fn captured_stderr(&self) -> Vec<Vec<u8>> {
if self.config.stderr_mode == StderrMode::Stdout {
return Vec::new();
}
self.shared
.queues
.lock()
.expect("queue mutex poisoned")
.stderr_history
.clone()
.into_iter()
.collect()
}
fn captured_stderr_raw(&self) -> Vec<u8> {
if self.config.stderr_mode == StderrMode::Stdout {
return Vec::new();
}
let guard = self.shared.queues.lock().expect("queue mutex poisoned");
guard.stderr_raw.iter().flatten().copied().collect()
}
pub fn captured_combined(&self) -> Vec<StreamEvent> {
self.shared
.queues
.lock()
.expect("queue mutex poisoned")
.combined_history
.clone()
.into_iter()
.collect()
}
pub fn captured_stream_bytes(&self, stream: StreamKind) -> usize {
if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
return 0;
}
let guard = self.shared.queues.lock().expect("queue mutex poisoned");
match stream {
StreamKind::Stdout => guard.stdout_history_bytes,
StreamKind::Stderr => guard.stderr_history_bytes,
}
}
pub fn captured_combined_bytes(&self) -> usize {
self.shared
.queues
.lock()
.expect("queue mutex poisoned")
.combined_history_bytes
}
pub fn clear_captured_stream(&self, stream: StreamKind) -> usize {
if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
return 0;
}
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
match stream {
StreamKind::Stdout => {
let released = guard.stdout_history_bytes;
guard.stdout_history.clear();
guard.stdout_raw.clear();
guard.stdout_raw_bytes = 0;
guard.stdout_history_bytes = 0;
released
}
StreamKind::Stderr => {
let released = guard.stderr_history_bytes;
guard.stderr_history.clear();
guard.stderr_raw.clear();
guard.stderr_raw_bytes = 0;
guard.stderr_history_bytes = 0;
released
}
}
}
pub fn clear_captured_combined(&self) -> usize {
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
let released = guard.combined_history_bytes;
guard.combined_history.clear();
guard.combined_history_bytes = 0;
released
}
fn build_command(&self) -> Command {
let command_override = self
.command_override
.lock()
.expect("command override mutex poisoned")
.take();
let mut command = match command_override {
Some(command) => command,
None => {
let mut command = match &self.config.command {
CommandSpec::Shell(command) => shell_command(command),
CommandSpec::Argv(argv) => {
let mut command = Command::new(&argv[0]);
if argv.len() > 1 {
command.args(&argv[1..]);
}
command
}
};
if let Some(cwd) = &self.config.cwd {
command.current_dir(cwd);
}
if let Some(env) = &self.config.env {
command.env_clear();
command.envs(env.iter().map(|(k, v)| (k, v)));
}
command
}
};
let platform_config =
running_process_platform_internal::platform::process::ProcessCommandConfig {
creation_flags: self.config.creationflags,
create_process_group: self.config.create_process_group,
nice: self.config.nice,
address_space_limit_bytes: self.config.address_space_limit_bytes,
};
let configured = if self.kill_when_owner_dies {
running_process_platform_internal::platform::process::
configure_process_command_for_bounded_owner_death(&mut command, platform_config)
} else {
running_process_platform_internal::platform::process::configure_process_command(
&mut command,
platform_config,
)
};
configured.expect("platform command configuration must be valid");
command
}
fn spawn_reader<R>(
&self,
pipe: R,
source_stream: StreamKind,
visible_stream: StreamKind,
on_pipe_done: Box<dyn FnOnce() + Send>,
) where
R: Read + Send + 'static,
{
let shared = Arc::clone(&self.shared);
shared.active_capture_readers.fetch_add(1, Ordering::AcqRel);
thread::spawn(move || {
let mut reader = pipe;
let mut chunk = vec![0_u8; 65536];
let mut pending = Vec::new();
loop {
match reader.read(&mut chunk) {
Ok(0) => break,
Ok(n) => {
if append_raw(&shared, visible_stream, &chunk[..n]) {
let lines = feed_chunk(&mut pending, &chunk[..n]);
emit_lines(&shared, visible_stream, lines);
} else {
pending.clear();
}
}
Err(_) => break,
}
}
if !pending.is_empty() && !shared.capture_overflowed.load(Ordering::Acquire) {
emit_lines(&shared, visible_stream, vec![std::mem::take(&mut pending)]);
}
on_pipe_done();
drop(reader);
let mut guard = shared.queues.lock().expect("queue mutex poisoned");
match source_stream {
StreamKind::Stdout => guard.stdout_closed = true,
StreamKind::Stderr => guard.stderr_closed = true,
}
shared.active_capture_readers.fetch_sub(1, Ordering::AcqRel);
shared.condvar.notify_all();
});
}
fn pipe_done_callback(&self, stream: StreamKind) -> Box<dyn FnOnce() + Send> {
let cancellation = Arc::clone(&self.capture_cancellation);
Box::new(move || {
let stream = match stream {
StreamKind::Stdout => {
running_process_platform_internal::platform::process::CaptureStream::Stdout
}
StreamKind::Stderr => {
running_process_platform_internal::platform::process::CaptureStream::Stderr
}
};
running_process_platform_internal::platform::process::capture_reader_done(
&cancellation,
stream,
);
})
}
fn cancel_capture_io(&self) {
crate::rp_rust_debug_scope!("running_process::NativeProcess::cancel_capture_io");
running_process_platform_internal::platform::process::cancel_capture_reader(
&self.capture_cancellation,
);
}
fn set_returncode(&self, code: i32) {
self.shared.returncode.store(code as i64, Ordering::Release);
self.shared.condvar.notify_all();
}
fn finish_capture_drain(&self) {
self.finish_capture_drain_with_deadline(kill_drain_deadline());
}
fn finish_capture_drain_with_deadline(&self, deadline: Instant) {
let drained = self.wait_for_capture_completion_with_deadline_impl(deadline);
if !drained {
self.cancel_capture_io();
}
}
fn wait_for_capture_completion_with_deadline_impl(&self, deadline: Instant) -> bool {
crate::rp_rust_debug_scope!(
"running_process::NativeProcess::wait_for_capture_completion_with_deadline"
);
if !self.config.capture {
return true;
}
finalize_capture_completion(&self.shared, deadline)
}
fn wait_for_capture_readers_with_deadline(&self, deadline: Instant) -> bool {
let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
while self.shared.active_capture_readers.load(Ordering::Acquire) != 0 {
let now = Instant::now();
if now >= deadline {
return false;
}
let (next_guard, result) = self
.shared
.condvar
.wait_timeout(guard, deadline - now)
.expect("queue mutex poisoned");
guard = next_guard;
if result.timed_out() && self.shared.active_capture_readers.load(Ordering::Acquire) != 0
{
return false;
}
}
true
}
}
fn finalize_capture_completion(shared: &SharedState, deadline: Instant) -> bool {
let mut guard = shared.queues.lock().expect("queue mutex poisoned");
while !(guard.stdout_closed && guard.stderr_closed) {
let now = Instant::now();
if now >= deadline {
guard.stdout_closed = true;
guard.stderr_closed = true;
shared.condvar.notify_all();
return false;
}
let (next_guard, result) = shared
.condvar
.wait_timeout(guard, deadline - now)
.expect("queue mutex poisoned");
guard = next_guard;
if result.timed_out() && !(guard.stdout_closed && guard.stderr_closed) {
guard.stdout_closed = true;
guard.stderr_closed = true;
shared.condvar.notify_all();
return false;
}
}
true
}
fn emit_lines(shared: &Arc<SharedState>, stream: StreamKind, lines: Vec<Vec<u8>>) {
if lines.is_empty() || shared.capture_overflowed.load(Ordering::Acquire) {
return;
}
let mut guard = shared.queues.lock().expect("queue mutex poisoned");
if shared.capture_overflowed.load(Ordering::Acquire) {
return;
}
for line in lines {
let line_len = line.len();
match stream {
StreamKind::Stdout => {
guard.stdout_history_bytes += line_len;
guard.stdout_history.push_back(line.clone());
guard.stdout_queue.push_back(line.clone());
}
StreamKind::Stderr => {
guard.stderr_history_bytes += line_len;
guard.stderr_history.push_back(line.clone());
guard.stderr_queue.push_back(line.clone());
}
}
let event = StreamEvent { stream, line };
guard.combined_history_bytes += line_len;
guard.combined_history.push_back(event.clone());
guard.combined_queue.push_back(event);
}
shared.condvar.notify_all();
}
fn append_raw(shared: &Arc<SharedState>, stream: StreamKind, chunk: &[u8]) -> bool {
if chunk.is_empty() {
return true;
}
let mut guard = shared.queues.lock().expect("queue mutex poisoned");
let accepted = match shared.capture_limit {
Some(limit) => {
let retained = guard
.stdout_raw_bytes
.saturating_add(guard.stderr_raw_bytes);
chunk.len().min(limit.saturating_sub(retained))
}
None => chunk.len(),
};
if accepted != 0 {
let accepted_chunk = chunk[..accepted].to_vec();
match stream {
StreamKind::Stdout => {
guard.stdout_raw_bytes += accepted;
guard.stdout_raw.push_back(accepted_chunk);
}
StreamKind::Stderr => {
guard.stderr_raw_bytes += accepted;
guard.stderr_raw.push_back(accepted_chunk);
}
}
}
if accepted != chunk.len() {
shared.capture_overflowed.store(true, Ordering::Release);
false
} else {
shared.condvar.notify_all();
true
}
}
mod bounded;
pub use bounded::{
run_command, run_command_bounded, run_std_command_bounded,
run_std_command_bounded_with_options, BoundedRunOptions,
};
pub(crate) fn shell_command(command: &str) -> Command {
running_process_platform_internal::platform::process::compat_shell_command(command)
}
#[cfg(test)]
mod tests;