1#![cfg_attr(
11 dylint_lib = "running_process_env_literal",
12 deny(running_process_env_direct)
13)]
14
15use std::collections::VecDeque;
16use std::io::Read;
17use std::process::{ChildStdin, Command, Stdio};
18use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
19use std::sync::{Arc, Condvar, Mutex};
20use std::thread;
21use std::time::{Duration, Instant};
22
23use crate::observer::{ObserverEmitter, ProcessWatchEmitter};
24
25pub use running_process_platform_internal::foreground;
27pub(crate) use running_process_platform_internal::platform;
28
29mod actor_runtime;
30pub mod ape;
31#[cfg(feature = "async-process")]
32mod async_process;
33#[cfg(feature = "async-process")]
34mod blocking_island;
35mod child_actor;
36#[cfg(feature = "async-process")]
37pub use blocking_island::dispatch_blocking as blocking_island_dispatch;
38pub mod console_detect;
39pub mod containment;
40mod descendant_monitor;
41pub mod env_vars;
42pub mod environment;
43mod helpers;
44#[cfg(feature = "async-process")]
45mod process_runtime;
46#[cfg(feature = "window-icon")]
47pub mod window_icon;
48#[cfg(feature = "daemon-registration")]
59pub mod daemon_registration;
60#[cfg(feature = "daemon-registration")]
62pub mod daemon_registration_compat;
63#[cfg(feature = "daemon-registration-v2")]
70pub mod daemon_registration_v2;
71#[cfg(feature = "daemon-registration-v2")]
73pub mod daemon_registration_v2_compat;
74#[cfg(feature = "frame-v1-codec")]
80pub mod daemon_frame_v1;
81#[cfg(any(feature = "daemon-registration", feature = "daemon-registration-v2"))]
82pub(crate) mod daemon_registration_common;
83#[cfg(feature = "frame-v1-codec")]
89pub mod frame_v1;
90#[cfg(any(feature = "backend-identity", feature = "daemon-registration"))]
94#[path = "broker/host_identity.rs"]
95pub(crate) mod daemon_host_identity;
96pub mod observer;
97#[cfg(feature = "originator-scan")]
98pub mod originator;
99pub mod output_log;
100#[cfg(feature = "client")]
104pub mod proto {
106 pub use running_process_protocol::daemon;
108}
109
110#[cfg(feature = "client")]
111pub mod client;
112
113#[cfg(feature = "client")]
119pub mod broker;
120#[cfg(all(feature = "backend-identity", not(feature = "client")))]
124#[allow(dead_code, unused_imports)]
125mod broker;
126
127#[cfg(feature = "backend-identity")]
134pub mod backend_identity;
135
136#[cfg(feature = "client")]
140pub mod content_hash;
141#[cfg(all(feature = "backend-identity", not(feature = "client")))]
142mod content_hash;
143
144#[cfg(feature = "probe")]
147pub mod probe;
148
149#[cfg(feature = "client")]
155pub mod maintenance;
156
157#[cfg(feature = "client")]
158pub mod cleanup;
159
160#[cfg(feature = "client")]
164pub mod boot_autostart;
165
166#[cfg(feature = "client")]
171pub mod runpm_config;
172
173#[cfg(feature = "test-support")]
178pub mod test_support;
179
180#[cfg(all(feature = "telemetry", not(feature = "daemon")))]
189#[path = "daemon/telemetry.rs"]
190pub mod telemetry;
191
192#[cfg(all(feature = "telemetry", feature = "daemon"))]
193pub use daemon::telemetry;
194
195#[cfg(all(feature = "telemetry", feature = "daemon"))]
205const _: fn(crate::telemetry::TeeHandle) -> daemon::telemetry::TeeHandle = |handle| handle;
206
207#[cfg(feature = "daemon")]
210pub mod daemon;
212#[cfg(feature = "independent-spawn")]
217pub mod independent_spawn;
218pub mod process_tree;
219#[cfg(feature = "pty")]
220pub mod pty;
222mod public_symbols;
223mod rust_debug;
224pub mod spawn;
225mod spawn_contract;
226pub use spawn_contract::{IndependentBackend, SpawnLifetime, SpawnMode, SpawnOptions};
227#[cfg(feature = "independent-spawn")]
228mod spawn_dispatch;
229#[cfg(feature = "independent-spawn")]
230pub use spawn_dispatch::{spawn_with_options, SpawnExit, SpawnHandle};
231pub mod systemd_killmode;
232#[cfg(feature = "terminal-graphics")]
233pub mod terminal_graphics;
234mod types;
235#[cfg(unix)]
236mod unix;
237mod windows;
238
239#[cfg(feature = "async-process")]
240pub use async_process::{
241 AsyncCapturedOutput, AsyncProcess, AsyncProcessBuilder, AsyncProcessSession,
242 AsyncProcessSessionChunk, AsyncProcessSessionControl, AsyncProcessSessionEvent,
243 AsyncProcessSessionOptions, AsyncProcessSessionOutput, AsyncStdio, ProcessTreeKill,
244};
245pub use console_detect::{monitor_console_windows, ConsoleWindowInfo};
246pub use containment::{ContainedProcessGroup, ORIGINATOR_ENV_VAR};
247#[cfg(feature = "client")]
249pub use content_hash::{blake3_file, daemon_identity_stamp, daemon_identity_stamp_env};
250pub use observer::{
251 CapabilitySupport, CaptureSource, CategoryCapability, DumpResult, EventCategory,
252 ObservationGrade, ObservationPolicy, ObserverCapabilities, ObserverConfig, ObserverEvent,
253 ObserverEventKind, ObserverSubscriber, ProcessEvent, ProcessEventKind, ProcessIdentity,
254 ProcessObservation, ProcessObservationCapabilities, ProcessObservationError, ProcessWatch,
255 ProcessWatchConfigurationError, ProcessWatchCursor, ProcessWatchGap, ProcessWatchLoss,
256 ProcessWatchMatch, ProcessWatchRead, ProcessWatchSubscriber, StackCapture, StackDump,
257};
258#[cfg(feature = "originator-scan")]
259pub use originator::{
260 find_declared_daemon_pids, find_processes_by_originator, OriginatorProcessInfo,
261};
262pub use output_log::{
263 CursorRead, OutputCursor, OutputLog, OutputRecord, SharedOutputCursor, SharedOutputLog,
264};
265#[doc(hidden)]
268pub use running_process_platform_internal::platform::executable as platform_executable;
269#[cfg(target_os = "linux")]
270pub use running_process_platform_internal::platform::process::current_executable_build_id;
271pub use running_process_platform_internal::platform::process::{
273 ProcessInspectError, ProcessInspectErrorKind,
274};
275pub use running_process_platform_internal::process_executable_path;
277pub use running_process_platform_internal::process_same_executable_path;
279pub use running_process_platform_internal::ProcessLiveness;
281pub use rust_debug::{render_rust_debug_traces, RustDebugScopeGuard};
282pub use spawn::{
283 spawn, spawn_daemon, spawn_daemon_breaking_away_from_job,
284 spawn_daemon_breaking_away_with_env_policy, spawn_daemon_with_clear_env,
285 spawn_daemon_with_env_policy, spawn_daemon_with_environment,
286 spawn_daemon_with_explicit_environment, spawn_daemon_with_stdio,
287 spawn_daemon_with_stdio_and_env_policy, spawn_with_env_policy, spawn_with_environment,
288 spawn_with_explicit_environment, DaemonChild, DaemonStdio, DaemonStdioSource,
289 EnvironmentPolicy, SpawnStdio, SpawnedChild, SpawnedChildControl, StdioSource, SyncEnvironment,
290 DAEMON_MARKER_ENV_VAR,
291};
292#[cfg(feature = "client-async")]
293pub use spawn::{spawn_tokio, TokioSpawnOptions};
294#[cfg(feature = "terminal-graphics")]
295pub use terminal_graphics::{
296 current_terminal_capabilities, current_terminal_capabilities_with_timeout,
297 detect_terminal_capabilities, CapabilityStatus, EvidenceStrength, GraphicsCapability,
298 GraphicsProtocol, TerminalCapabilities, TerminalCapabilityInput, TerminalGraphicsCapabilities,
299 TerminalProbeEvidence,
300};
301pub use types::{
302 CommandSpec, ProcessConfig, ProcessError, ReadStatus, RunOutput, StderrMode, StdinMode,
303 StreamEvent, StreamKind,
304};
305#[cfg(feature = "window-icon")]
306pub use window_icon::{
307 host_icon_support, icon_support, set_host_icon, set_icon, IconError, IconScope, IconSource,
308 IconSupport, StockIcon,
309};
310
311pub(crate) use helpers::child_try_wait_error_is_retryable;
312#[cfg(test)]
313pub(crate) use helpers::exit_code;
314pub(crate) use helpers::{feed_chunk, kill_drain_deadline, log_spawned_child_pid};
315pub use running_process_platform_internal::exit_code as native_exit_code;
317pub use running_process_platform_internal::ProcessPriority;
318#[cfg(feature = "async-process")]
319pub use running_process_platform_internal::SpawnAdmission;
320#[cfg(unix)]
321pub use unix::{unix_set_priority, unix_signal_process, unix_signal_process_group, UnixSignal};
322pub(crate) use windows::{assign_child_to_windows_kill_on_close_job_impl, WindowsJobHandle};
323
324#[macro_export]
325macro_rules! rp_rust_debug_scope {
327 ($label:expr) => {
328 let _running_process_rust_debug_scope =
329 $crate::RustDebugScopeGuard::enter($label, file!(), line!());
330 };
331}
332
333#[derive(Default)]
334struct QueueState {
335 stdout_queue: VecDeque<Vec<u8>>,
336 stderr_queue: VecDeque<Vec<u8>>,
337 combined_queue: VecDeque<StreamEvent>,
338 stdout_history: VecDeque<Vec<u8>>,
339 stderr_history: VecDeque<Vec<u8>>,
340 combined_history: VecDeque<StreamEvent>,
341 stdout_raw: VecDeque<Vec<u8>>,
345 stderr_raw: VecDeque<Vec<u8>>,
346 stdout_raw_bytes: usize,
347 stderr_raw_bytes: usize,
348 stdout_history_bytes: usize,
349 stderr_history_bytes: usize,
350 combined_history_bytes: usize,
351 stdout_closed: bool,
352 stderr_closed: bool,
353}
354
355const RETURNCODE_NOT_SET: i64 = i64::MIN;
357
358struct SharedState {
359 queues: Mutex<QueueState>,
360 condvar: Condvar,
361 capture_limit: Option<usize>,
362 capture_overflowed: AtomicBool,
363 active_capture_readers: std::sync::atomic::AtomicUsize,
364 returncode: AtomicI64,
367 observer: Option<ObserverEmitter>,
372 observer_exit_emitted: AtomicBool,
375 exit_code: tokio::sync::watch::Sender<Option<i32>>,
378}
379
380type ChildState = running_process_platform_internal::platform::process::PlatformStdChild;
385
386#[cfg(test)]
387#[derive(Debug, Eq, PartialEq)]
388enum CapturePollAction {
389 Wait,
390 Read,
391 Cancel,
392}
393
394#[cfg(test)]
395fn capture_poll_action(capture_revents: i16, wake_revents: i16) -> CapturePollAction {
396 if wake_revents != 0 {
397 CapturePollAction::Cancel
398 } else if capture_revents != 0 {
399 CapturePollAction::Read
400 } else {
401 CapturePollAction::Wait
402 }
403}
404
405fn cleanup_child_after_start_error(child: ChildState) {
406 child.discard_after_start_error();
407}
408
409impl SharedState {
410 #[cfg(test)]
411 fn new(capture: bool) -> Self {
412 Self::with_observer_and_limit(capture, None, None)
413 }
414
415 fn with_observer_and_limit(
416 capture: bool,
417 observer: Option<ObserverEmitter>,
418 capture_limit: Option<usize>,
419 ) -> Self {
420 let queues = QueueState {
421 stdout_closed: !capture,
422 stderr_closed: !capture,
423 ..QueueState::default()
424 };
425 Self {
426 queues: Mutex::new(queues),
427 condvar: Condvar::new(),
428 capture_limit,
429 capture_overflowed: AtomicBool::new(false),
430 active_capture_readers: std::sync::atomic::AtomicUsize::new(0),
431 returncode: AtomicI64::new(RETURNCODE_NOT_SET),
432 observer,
433 observer_exit_emitted: AtomicBool::new(false),
434 exit_code: tokio::sync::watch::Sender::new(None),
435 }
436 }
437
438 fn record_exit(&self, code: i32) {
441 self.returncode.store(code as i64, Ordering::Release);
442 self.exit_code.send_replace(Some(code));
443 self.condvar.notify_all();
444 }
445
446 fn emit_exited(&self, pid: u32, exit_code: i32) {
449 let Some(emitter) = self.observer.as_ref() else {
450 return;
451 };
452 if self
453 .observer_exit_emitted
454 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
455 .is_ok()
456 {
457 emitter.emit_exited(pid, exit_code);
458 }
459 }
460}
461
462pub struct NativeProcess {
469 config: ProcessConfig,
470 command_override: Mutex<Option<Command>>,
471 child: std::sync::OnceLock<child_actor::ChildHandle>,
473 start_gate: Mutex<()>,
475 stdin: Mutex<Option<ChildStdin>>,
476 shared: Arc<SharedState>,
477 process_watch: Option<Arc<ProcessWatchEmitter>>,
478 kill_when_owner_dies: bool,
482 #[cfg(test)]
483 stdin_write_active: AtomicBool,
484 capture_cancellation:
485 Arc<running_process_platform_internal::platform::process::CaptureCancellation>,
486}
487
488impl NativeProcess {
489 pub fn new(config: ProcessConfig) -> Self {
495 Self::new_with_options(config, None, None, None, None)
496 }
497
498 pub fn with_observer(
511 config: ProcessConfig,
512 observer: crate::observer::ObserverConfig,
513 ) -> (Self, ObserverSubscriber) {
514 let (emitter, subscriber) = ObserverEmitter::new(observer);
515 let process = Self::new_with_options(config, Some(emitter), None, None, None);
516 (process, subscriber)
517 }
518
519 pub fn with_observer_and_command(
534 command: Command,
535 config: ProcessConfig,
536 observer: crate::observer::ObserverConfig,
537 ) -> (Self, ObserverSubscriber) {
538 let (emitter, subscriber) = ObserverEmitter::new(observer);
539 let process = Self::new_with_options(config, Some(emitter), None, Some(command), None);
540 (process, subscriber)
541 }
542
543 pub fn with_process_watches(
546 config: ProcessConfig,
547 watches: Vec<ProcessWatch>,
548 policy: ObservationPolicy,
549 ) -> Result<(Self, ProcessWatchSubscriber), ProcessObservationError> {
550 let (emitter, subscriber) = ProcessWatchEmitter::new(watches, policy)?;
551 let process = Self::new_with_options(config, None, None, None, Some(emitter));
552 Ok((process, subscriber))
553 }
554
555 pub fn process_observation_capabilities() -> ProcessObservationCapabilities {
557 ProcessObservationCapabilities::current()
558 }
559
560 fn new_with_capture_limit(config: ProcessConfig, capture_limit: usize) -> Self {
561 Self::new_with_options(config, None, Some(capture_limit), None, None)
562 }
563
564 fn new_with_command_capture_limit(
565 command: Command,
566 config: ProcessConfig,
567 capture_limit: usize,
568 kill_when_owner_dies: bool,
569 ) -> Self {
570 let mut process =
571 Self::new_with_options(config, None, Some(capture_limit), Some(command), None);
572 process.kill_when_owner_dies = kill_when_owner_dies;
573 process
574 }
575
576 fn new_with_options(
577 config: ProcessConfig,
578 observer: Option<ObserverEmitter>,
579 capture_limit: Option<usize>,
580 command_override: Option<Command>,
581 process_watch: Option<Arc<ProcessWatchEmitter>>,
582 ) -> Self {
583 let shared = SharedState::with_observer_and_limit(config.capture, observer, capture_limit);
584 Self {
585 shared: Arc::new(shared),
586 process_watch,
587 command_override: Mutex::new(command_override),
588 child: std::sync::OnceLock::new(),
589 start_gate: Mutex::new(()),
590 stdin: Mutex::new(None),
591 kill_when_owner_dies: false,
592 #[cfg(test)]
593 stdin_write_active: AtomicBool::new(false),
594 config,
595 capture_cancellation: Arc::new(Default::default()),
596 }
597 }
598
599 #[inline(never)]
601 pub fn start(&self) -> Result<(), ProcessError> {
606 public_symbols::rp_native_process_start_public(self)
607 }
608
609 fn start_impl(&self) -> Result<(), ProcessError> {
610 crate::rp_rust_debug_scope!("running_process::NativeProcess::start");
611 let _gate = self.start_gate.lock().expect("start gate poisoned");
612 if self.child.get().is_some() {
613 return Err(ProcessError::AlreadyStarted);
614 }
615
616 let mut command = self.build_command();
617 let exact_trace = self
618 .process_watch
619 .as_ref()
620 .is_some_and(|watch| watch.uses_exact_trace());
621 match self.config.stdin_mode {
622 StdinMode::Inherit => {}
623 StdinMode::Piped => {
624 command.stdin(Stdio::piped());
625 }
626 StdinMode::Null => {
627 command.stdin(Stdio::null());
628 }
629 }
630 if self.config.capture {
631 command.stdout(Stdio::piped());
632 command.stderr(Stdio::piped());
633 }
634
635 let mut child = if exact_trace {
636 let event_watch = Arc::clone(self.process_watch.as_ref().expect("exact watch checked"));
637 let completion_watch = Arc::clone(&event_watch);
638 match running_process_platform_internal::platform::process::start_exact_trace(
639 command,
640 Box::new(move |event| event_watch.emit_exact(event)),
641 Box::new(move || completion_watch.close()),
642 ) {
643 Ok(child) => ChildState::from_exact_trace(child, self.config.create_process_group),
644 Err(error) => {
645 if let Some(watch) = self.process_watch.as_ref() {
646 watch.close();
647 }
648 return Err(ProcessError::Spawn(error));
649 }
650 }
651 } else {
652 ChildState::from_std(
653 running_process_platform_internal::platform::ape::spawn_std(
654 &mut command,
655 |command| command.spawn(),
656 )
657 .map_err(ProcessError::Spawn)?,
658 self.config.create_process_group,
659 )
660 };
661 log_spawned_child_pid(child.id()).map_err(ProcessError::Spawn)?;
662 if let Some(emitter) = self.shared.observer.as_ref() {
665 emitter.emit_started(child.id());
666 }
667 let job_result = child.std_child().map(|standard_child| {
675 let descendant_sink = self
676 .shared
677 .observer
678 .as_ref()
679 .and_then(|e| e.descendant_sink());
680 public_symbols::rp_assign_child_to_windows_kill_on_close_job_with_observer_public(
681 standard_child,
682 descendant_sink,
683 self.process_watch.clone(),
684 standard_child.id(),
685 self.config.address_space_limit_bytes,
686 )
687 });
688 match job_result {
689 None => {}
690 Some(Ok(job)) => child.attach_job(job),
691 Some(Err(error)) if error.kind() == std::io::ErrorKind::Unsupported => {}
692 Some(Err(error)) => {
693 if let Some(watch) = self.process_watch.as_ref() {
694 watch.close();
695 }
696 cleanup_child_after_start_error(child);
697 return Err(ProcessError::Spawn(error));
698 }
699 }
700 if !exact_trace {
701 descendant_monitor::start(
702 child.id(),
703 self.shared.observer.as_ref(),
704 self.process_watch.as_ref(),
705 );
706 }
707 if self.config.capture {
708 let readers = match child.prepare_capture(&self.capture_cancellation) {
709 Ok(readers) => readers,
710 Err(error) => {
711 cleanup_child_after_start_error(child);
712 return Err(ProcessError::Spawn(error));
713 }
714 };
715 let (stdout, stderr) = (readers.stdout, readers.stderr);
716 self.spawn_reader(
717 stdout,
718 StreamKind::Stdout,
719 StreamKind::Stdout,
720 self.pipe_done_callback(StreamKind::Stdout),
721 );
722 self.spawn_reader(
723 stderr,
724 StreamKind::Stderr,
725 match self.config.stderr_mode {
726 StderrMode::Stdout => StreamKind::Stdout,
727 StderrMode::Pipe => StreamKind::Stderr,
728 },
729 self.pipe_done_callback(StreamKind::Stderr),
730 );
731 }
732 *self.stdin.lock().expect("stdin mutex poisoned") = child.take_stdin();
733 let handle = child_actor::spawn(
734 child,
735 Arc::clone(&self.shared),
736 self.config.capture,
737 Arc::clone(&self.capture_cancellation),
738 );
739 let _ = self.child.set(handle);
740 Ok(())
741 }
742
743 pub fn write_stdin(&self, data: &[u8]) -> Result<(), ProcessError> {
745 if self.child.get().is_none() {
746 return Err(ProcessError::NotRunning);
747 }
748 let mut guard = self.stdin.lock().expect("stdin mutex poisoned");
749 let stdin = guard.as_mut().ok_or(ProcessError::StdinUnavailable)?;
750 use std::io::Write;
751 #[cfg(test)]
752 self.stdin_write_active.store(true, Ordering::Release);
753 let write_result = stdin.write_all(data);
754 #[cfg(test)]
755 self.stdin_write_active.store(false, Ordering::Release);
756 write_result.map_err(ProcessError::Io)?;
757 stdin.flush().map_err(ProcessError::Io)?;
758 drop(guard.take());
759 Ok(())
760 }
761
762 pub fn write_stdin_streaming(&self, data: &[u8]) -> Result<(), ProcessError> {
767 if self.child.get().is_none() {
768 return Err(ProcessError::NotRunning);
769 }
770 let mut guard = self.stdin.lock().expect("stdin mutex poisoned");
771 let stdin = guard.as_mut().ok_or(ProcessError::StdinUnavailable)?;
772 use std::io::Write;
773 #[cfg(test)]
774 self.stdin_write_active.store(true, Ordering::Release);
775 let write_result = stdin.write_all(data);
776 #[cfg(test)]
777 self.stdin_write_active.store(false, Ordering::Release);
778 write_result.map_err(ProcessError::Io)?;
779 stdin.flush().map_err(ProcessError::Io)?;
780 Ok(())
781 }
782
783 pub fn close_stdin(&self) -> Result<(), ProcessError> {
786 if self.child.get().is_none() {
787 return Err(ProcessError::NotRunning);
788 }
789 drop(self.stdin.lock().expect("stdin mutex poisoned").take());
790 Ok(())
791 }
792
793 pub fn poll(&self) -> Result<Option<i32>, ProcessError> {
797 if let Some(code) = self.returncode() {
799 return Ok(Some(code));
800 }
801 let Some(child) = self.child.get() else {
802 return Ok(self.returncode());
803 };
804 child.try_wait().map_err(ProcessError::Io)
807 }
808
809 #[inline(never)]
811 pub fn wait(&self, timeout: Option<Duration>) -> Result<i32, ProcessError> {
816 public_symbols::rp_native_process_wait_public(self, timeout)
817 }
818
819 fn wait_impl(&self, timeout: Option<Duration>) -> Result<i32, ProcessError> {
820 crate::rp_rust_debug_scope!("running_process::NativeProcess::wait");
821 if self.child.get().is_none() {
822 return self.returncode().ok_or(ProcessError::NotRunning);
823 }
824 if let Some(code) = self.returncode() {
826 self.finish_capture_drain();
827 return Ok(code);
828 }
829 let mut timeout = timeout;
837 if let Some(limit) = timeout {
838 let fine = limit.min(SHORT_WAIT_DIRECT_POLL);
839 let deadline = Instant::now() + fine;
840 loop {
841 if let Some(child) = self.child.get() {
844 let _ = child.try_wait();
845 }
846 if let Some(code) = self.returncode() {
847 self.finish_capture_drain();
848 return Ok(code);
849 }
850 let now = Instant::now();
851 if now >= deadline {
852 break;
853 }
854 thread::sleep((deadline - now).min(Duration::from_millis(1)));
855 }
856 if let Some(code) = self.returncode() {
857 self.finish_capture_drain();
858 return Ok(code);
859 }
860 let remaining = limit.saturating_sub(fine);
861 if remaining.is_zero() {
862 return Err(ProcessError::Timeout);
863 }
864 timeout = Some(remaining);
865 }
866 let mut exit = self.shared.exit_code.subscribe();
870 let outcome = actor_runtime::block_on_anywhere(async move {
871 let exited = async move {
872 exit.wait_for(Option::is_some)
873 .await
874 .ok()
875 .and_then(|code| *code)
876 };
877 match timeout {
878 Some(limit) => tokio::time::timeout(limit, exited).await.ok(),
879 None => Some(exited.await),
880 }
881 });
882 match outcome {
883 None => Err(ProcessError::Timeout),
884 Some(None) => Err(ProcessError::NotRunning),
887 Some(Some(code)) => {
888 self.finish_capture_drain();
889 Ok(code)
890 }
891 }
892 }
893
894 #[inline(never)]
896 pub fn kill(&self) -> Result<(), ProcessError> {
898 public_symbols::rp_native_process_kill_public(self)
899 }
900
901 fn kill_impl(&self) -> Result<(), ProcessError> {
902 crate::rp_rust_debug_scope!("running_process::NativeProcess::kill");
903 let deadline = kill_drain_deadline();
904 let child = self.child.get().ok_or(ProcessError::NotRunning)?;
905 let already_reaped = match child.kill().map_err(ProcessError::Io)? {
908 child_actor::KillOutcome::AlreadyExited(code) => Some(code),
909 child_actor::KillOutcome::Signalled => None,
910 };
911
912 self.cancel_capture_io();
918 if already_reaped.is_none() {
922 self.await_exit_until(deadline);
923 }
924 public_symbols::rp_native_process_wait_for_capture_completion_with_deadline_public(
930 self, deadline,
931 );
932 Ok(())
933 }
934
935 fn await_exit_until(&self, deadline: Instant) -> Option<i32> {
939 let mut exit = self.shared.exit_code.subscribe();
940 let remaining = deadline.saturating_duration_since(Instant::now());
941 let published = actor_runtime::block_on_anywhere(async move {
942 tokio::time::timeout(remaining, exit.wait_for(Option::is_some))
943 .await
944 .ok()
945 .and_then(|seen| seen.ok().and_then(|code| *code))
946 });
947 published.or_else(|| self.child.get()?.try_wait().ok().flatten())
948 }
949
950 pub fn terminate(&self) -> Result<(), ProcessError> {
954 self.kill()
955 }
956
957 pub fn terminate_group_soft(&self) -> Result<(), ProcessError> {
969 if !self.config.create_process_group {
970 return Ok(());
972 }
973 let pid = self.pid().ok_or(ProcessError::NotRunning)?;
974 running_process_platform_internal::platform::process::soft_terminate_process_group(pid)
975 .map_err(ProcessError::Io)
976 }
977
978 #[inline(never)]
980 pub fn close(&self) -> Result<(), ProcessError> {
982 public_symbols::rp_native_process_close_public(self)
983 }
984
985 fn close_impl(&self) -> Result<(), ProcessError> {
986 crate::rp_rust_debug_scope!("running_process::NativeProcess::close");
987 if self.child.get().is_none() {
988 return Ok(());
989 }
990 if self.poll()?.is_none() {
991 self.kill()?;
992 } else {
993 self.finish_capture_drain();
994 }
995 if let Some(watch) = self.process_watch.as_ref() {
996 watch.close();
997 }
998 Ok(())
999 }
1000
1001 pub fn pid(&self) -> Option<u32> {
1003 self.child.get().map(child_actor::ChildHandle::pid)
1004 }
1005
1006 pub fn returncode(&self) -> Option<i32> {
1008 let v = self.shared.returncode.load(Ordering::Acquire);
1009 if v == RETURNCODE_NOT_SET {
1010 None
1011 } else {
1012 Some(v as i32)
1013 }
1014 }
1015
1016 pub fn has_pending_stream(&self, stream: StreamKind) -> bool {
1018 if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
1019 return false;
1020 }
1021 let guard = self.shared.queues.lock().expect("queue mutex poisoned");
1022 match stream {
1023 StreamKind::Stdout => !guard.stdout_queue.is_empty(),
1024 StreamKind::Stderr => !guard.stderr_queue.is_empty(),
1025 }
1026 }
1027
1028 pub fn has_pending_combined(&self) -> bool {
1030 let guard = self.shared.queues.lock().expect("queue mutex poisoned");
1031 !guard.combined_queue.is_empty()
1032 }
1033
1034 pub fn drain_stream(&self, stream: StreamKind) -> Vec<Vec<u8>> {
1036 if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
1037 return Vec::new();
1038 }
1039 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1040 let queue = match stream {
1041 StreamKind::Stdout => &mut guard.stdout_queue,
1042 StreamKind::Stderr => &mut guard.stderr_queue,
1043 };
1044 queue.drain(..).collect()
1045 }
1046
1047 pub fn drain_stream_raw(&self, stream: StreamKind) -> Vec<u8> {
1056 if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
1057 return Vec::new();
1058 }
1059 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1060 match stream {
1061 StreamKind::Stdout => {
1062 let mut output = Vec::with_capacity(guard.stdout_raw_bytes);
1063 for chunk in guard.stdout_raw.drain(..) {
1064 output.extend_from_slice(&chunk);
1065 }
1066 guard.stdout_raw_bytes = 0;
1067 output
1068 }
1069 StreamKind::Stderr => {
1070 let mut output = Vec::with_capacity(guard.stderr_raw_bytes);
1071 for chunk in guard.stderr_raw.drain(..) {
1072 output.extend_from_slice(&chunk);
1073 }
1074 guard.stderr_raw_bytes = 0;
1075 output
1076 }
1077 }
1078 }
1079
1080 pub fn drain_combined(&self) -> Vec<StreamEvent> {
1082 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1083 guard.combined_queue.drain(..).collect()
1084 }
1085
1086 pub fn read_stream(
1091 &self,
1092 stream: StreamKind,
1093 timeout: Option<Duration>,
1094 ) -> ReadStatus<Vec<u8>> {
1095 let deadline = timeout.map(|limit| Instant::now() + limit);
1096 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1097
1098 loop {
1099 if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
1100 return ReadStatus::Eof;
1101 }
1102
1103 let queue = match stream {
1104 StreamKind::Stdout => &mut guard.stdout_queue,
1105 StreamKind::Stderr => &mut guard.stderr_queue,
1106 };
1107 if let Some(line) = queue.pop_front() {
1108 return ReadStatus::Line(line);
1109 }
1110
1111 let closed = match stream {
1112 StreamKind::Stdout => {
1113 if self.config.stderr_mode == StderrMode::Stdout {
1114 guard.stdout_closed && guard.stderr_closed
1115 } else {
1116 guard.stdout_closed
1117 }
1118 }
1119 StreamKind::Stderr => guard.stderr_closed,
1120 };
1121 if closed {
1122 return ReadStatus::Eof;
1123 }
1124
1125 match deadline {
1126 Some(deadline) => {
1127 let now = Instant::now();
1128 if now >= deadline {
1129 return ReadStatus::Timeout;
1130 }
1131 let wait = deadline.saturating_duration_since(now);
1132 let result = self
1133 .shared
1134 .condvar
1135 .wait_timeout(guard, wait)
1136 .expect("queue mutex poisoned");
1137 guard = result.0;
1138 if result.1.timed_out() {
1139 return ReadStatus::Timeout;
1140 }
1141 }
1142 None => {
1143 guard = self
1144 .shared
1145 .condvar
1146 .wait(guard)
1147 .expect("queue mutex poisoned");
1148 }
1149 }
1150 }
1151 }
1152
1153 #[inline(never)]
1155 pub fn read_combined(&self, timeout: Option<Duration>) -> ReadStatus<StreamEvent> {
1157 public_symbols::rp_native_process_read_combined_public(self, timeout)
1158 }
1159
1160 fn read_combined_impl(&self, timeout: Option<Duration>) -> ReadStatus<StreamEvent> {
1161 crate::rp_rust_debug_scope!("running_process::NativeProcess::read_combined");
1162 let deadline = timeout.map(|limit| Instant::now() + limit);
1163 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1164
1165 loop {
1166 if let Some(event) = guard.combined_queue.pop_front() {
1167 return ReadStatus::Line(event);
1168 }
1169 if guard.stdout_closed && guard.stderr_closed {
1170 return ReadStatus::Eof;
1171 }
1172
1173 match deadline {
1174 Some(deadline) => {
1175 let now = Instant::now();
1176 if now >= deadline {
1177 return ReadStatus::Timeout;
1178 }
1179 let wait = deadline.saturating_duration_since(now);
1180 let result = self
1181 .shared
1182 .condvar
1183 .wait_timeout(guard, wait)
1184 .expect("queue mutex poisoned");
1185 guard = result.0;
1186 if result.1.timed_out() {
1187 return ReadStatus::Timeout;
1188 }
1189 }
1190 None => {
1191 guard = self
1192 .shared
1193 .condvar
1194 .wait(guard)
1195 .expect("queue mutex poisoned");
1196 }
1197 }
1198 }
1199 }
1200
1201 pub fn captured_stdout(&self) -> Vec<Vec<u8>> {
1203 self.shared
1204 .queues
1205 .lock()
1206 .expect("queue mutex poisoned")
1207 .stdout_history
1208 .clone()
1209 .into_iter()
1210 .collect()
1211 }
1212
1213 fn captured_stdout_raw(&self) -> Vec<u8> {
1214 let guard = self.shared.queues.lock().expect("queue mutex poisoned");
1215 guard.stdout_raw.iter().flatten().copied().collect()
1216 }
1217
1218 pub fn captured_stderr(&self) -> Vec<Vec<u8>> {
1220 if self.config.stderr_mode == StderrMode::Stdout {
1221 return Vec::new();
1222 }
1223 self.shared
1224 .queues
1225 .lock()
1226 .expect("queue mutex poisoned")
1227 .stderr_history
1228 .clone()
1229 .into_iter()
1230 .collect()
1231 }
1232
1233 fn captured_stderr_raw(&self) -> Vec<u8> {
1234 if self.config.stderr_mode == StderrMode::Stdout {
1235 return Vec::new();
1236 }
1237 let guard = self.shared.queues.lock().expect("queue mutex poisoned");
1238 guard.stderr_raw.iter().flatten().copied().collect()
1239 }
1240
1241 pub fn captured_combined(&self) -> Vec<StreamEvent> {
1243 self.shared
1244 .queues
1245 .lock()
1246 .expect("queue mutex poisoned")
1247 .combined_history
1248 .clone()
1249 .into_iter()
1250 .collect()
1251 }
1252
1253 pub fn captured_stream_bytes(&self, stream: StreamKind) -> usize {
1255 if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
1256 return 0;
1257 }
1258 let guard = self.shared.queues.lock().expect("queue mutex poisoned");
1259 match stream {
1260 StreamKind::Stdout => guard.stdout_history_bytes,
1261 StreamKind::Stderr => guard.stderr_history_bytes,
1262 }
1263 }
1264
1265 pub fn captured_combined_bytes(&self) -> usize {
1267 self.shared
1268 .queues
1269 .lock()
1270 .expect("queue mutex poisoned")
1271 .combined_history_bytes
1272 }
1273
1274 pub fn clear_captured_stream(&self, stream: StreamKind) -> usize {
1283 if stream == StreamKind::Stderr && self.config.stderr_mode == StderrMode::Stdout {
1284 return 0;
1285 }
1286 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1287 match stream {
1288 StreamKind::Stdout => {
1289 let released = guard.stdout_history_bytes;
1290 guard.stdout_history.clear();
1291 guard.stdout_raw.clear();
1292 guard.stdout_raw_bytes = 0;
1293 guard.stdout_history_bytes = 0;
1294 released
1295 }
1296 StreamKind::Stderr => {
1297 let released = guard.stderr_history_bytes;
1298 guard.stderr_history.clear();
1299 guard.stderr_raw.clear();
1300 guard.stderr_raw_bytes = 0;
1301 guard.stderr_history_bytes = 0;
1302 released
1303 }
1304 }
1305 }
1306
1307 pub fn clear_captured_combined(&self) -> usize {
1309 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1310 let released = guard.combined_history_bytes;
1311 guard.combined_history.clear();
1312 guard.combined_history_bytes = 0;
1313 released
1314 }
1315
1316 fn build_command(&self) -> Command {
1317 let command_override = self
1318 .command_override
1319 .lock()
1320 .expect("command override mutex poisoned")
1321 .take();
1322 let mut command = match command_override {
1323 Some(command) => command,
1324 None => {
1325 let mut ape_path = None;
1328 let mut command = match &self.config.command {
1329 CommandSpec::Shell(command) => shell_command(command),
1330 CommandSpec::Argv(argv) => {
1331 let options = platform::ape::ApeOptions::with_overrides(
1334 self.config.env.is_some(),
1335 self.config.env.iter().flatten().map(|(key, value)| {
1336 (std::ffi::OsStr::new(key), Some(std::ffi::OsStr::new(value)))
1337 }),
1338 );
1339 match platform::ape::plan_launch(
1340 std::ffi::OsStr::new(&argv[0]),
1341 self.config.cwd.as_deref(),
1342 &options,
1343 ) {
1344 Some(launch) => {
1345 ape_path = launch.child_path(options.path.as_deref());
1346 let mut command = Command::new(&launch.loader);
1347 command.args(launch.args(&argv[1..]));
1348 command
1349 }
1350 None => {
1351 let mut command = Command::new(&argv[0]);
1352 command.args(&argv[1..]);
1353 command
1354 }
1355 }
1356 }
1357 };
1358 if let Some(cwd) = &self.config.cwd {
1359 command.current_dir(cwd);
1360 }
1361 if let Some(env) = &self.config.env {
1362 command.env_clear();
1363 command.envs(env.iter().map(|(k, v)| (k, v)));
1364 }
1365 if let Some(path) = ape_path {
1366 command.env("PATH", path);
1367 }
1368 command
1369 }
1370 };
1371 let platform_config =
1372 running_process_platform_internal::platform::process::ProcessCommandConfig {
1373 creation_flags: self.config.creationflags,
1374 create_process_group: self.config.create_process_group,
1375 nice: self.config.nice,
1376 address_space_limit_bytes: self.config.address_space_limit_bytes,
1377 };
1378 let configured = if self.kill_when_owner_dies {
1379 running_process_platform_internal::platform::process::
1380 configure_process_command_for_bounded_owner_death(&mut command, platform_config)
1381 } else {
1382 running_process_platform_internal::platform::process::configure_process_command(
1383 &mut command,
1384 platform_config,
1385 )
1386 };
1387 configured.expect("platform command configuration must be valid");
1388 command
1389 }
1390
1391 fn spawn_reader<R>(
1392 &self,
1393 pipe: R,
1394 source_stream: StreamKind,
1395 visible_stream: StreamKind,
1396 on_pipe_done: Box<dyn FnOnce() + Send>,
1397 ) where
1398 R: Read + Send + 'static,
1399 {
1400 let shared = Arc::clone(&self.shared);
1401 shared.active_capture_readers.fetch_add(1, Ordering::AcqRel);
1402 thread::spawn(move || {
1403 let mut reader = pipe;
1404 let mut chunk = vec![0_u8; 65536];
1405 let mut pending = Vec::new();
1406
1407 loop {
1408 match reader.read(&mut chunk) {
1409 Ok(0) => break,
1410 Ok(n) => {
1411 if append_raw(&shared, visible_stream, &chunk[..n]) {
1412 let lines = feed_chunk(&mut pending, &chunk[..n]);
1413 emit_lines(&shared, visible_stream, lines);
1414 } else {
1415 pending.clear();
1416 }
1417 }
1418 Err(_) => break,
1419 }
1420 }
1421
1422 if !pending.is_empty() && !shared.capture_overflowed.load(Ordering::Acquire) {
1423 emit_lines(&shared, visible_stream, vec![std::mem::take(&mut pending)]);
1424 }
1425
1426 on_pipe_done();
1431 drop(reader);
1432
1433 let mut guard = shared.queues.lock().expect("queue mutex poisoned");
1434 match source_stream {
1435 StreamKind::Stdout => guard.stdout_closed = true,
1436 StreamKind::Stderr => guard.stderr_closed = true,
1437 }
1438 shared.active_capture_readers.fetch_sub(1, Ordering::AcqRel);
1439 shared.condvar.notify_all();
1440 });
1441 }
1442
1443 fn pipe_done_callback(&self, stream: StreamKind) -> Box<dyn FnOnce() + Send> {
1444 let cancellation = Arc::clone(&self.capture_cancellation);
1445 Box::new(move || {
1446 let stream = match stream {
1447 StreamKind::Stdout => {
1448 running_process_platform_internal::platform::process::CaptureStream::Stdout
1449 }
1450 StreamKind::Stderr => {
1451 running_process_platform_internal::platform::process::CaptureStream::Stderr
1452 }
1453 };
1454 running_process_platform_internal::platform::process::capture_reader_done(
1455 &cancellation,
1456 stream,
1457 );
1458 })
1459 }
1460
1461 fn cancel_capture_io(&self) {
1465 crate::rp_rust_debug_scope!("running_process::NativeProcess::cancel_capture_io");
1466 running_process_platform_internal::platform::process::cancel_capture_reader(
1467 &self.capture_cancellation,
1468 );
1469 }
1470
1471 #[cfg(test)]
1472 fn set_returncode(&self, code: i32) {
1473 self.shared.record_exit(code);
1474 }
1475
1476 fn finish_capture_drain(&self) {
1486 self.finish_capture_drain_with_deadline(kill_drain_deadline());
1487 }
1488
1489 fn finish_capture_drain_with_deadline(&self, deadline: Instant) {
1490 let drained = self.wait_for_capture_completion_with_deadline_impl(deadline);
1491 if !drained {
1492 self.cancel_capture_io();
1493 }
1494 }
1495
1496 fn wait_for_capture_completion_with_deadline_impl(&self, deadline: Instant) -> bool {
1499 crate::rp_rust_debug_scope!(
1500 "running_process::NativeProcess::wait_for_capture_completion_with_deadline"
1501 );
1502 if !self.config.capture {
1503 return true;
1504 }
1505 finalize_capture_completion(&self.shared, deadline)
1506 }
1507
1508 fn wait_for_capture_readers_with_deadline(&self, deadline: Instant) -> bool {
1509 let mut guard = self.shared.queues.lock().expect("queue mutex poisoned");
1510 while self.shared.active_capture_readers.load(Ordering::Acquire) != 0 {
1511 let now = Instant::now();
1512 if now >= deadline {
1513 return false;
1514 }
1515 let (next_guard, result) = self
1516 .shared
1517 .condvar
1518 .wait_timeout(guard, deadline - now)
1519 .expect("queue mutex poisoned");
1520 guard = next_guard;
1521 if result.timed_out() && self.shared.active_capture_readers.load(Ordering::Acquire) != 0
1522 {
1523 return false;
1524 }
1525 }
1526 true
1527 }
1528}
1529
1530const SHORT_WAIT_DIRECT_POLL: Duration = Duration::from_millis(50);
1534
1535pub(crate) async fn finalize_capture_completion_async(
1549 shared: &SharedState,
1550 deadline: Instant,
1551) -> bool {
1552 loop {
1553 {
1554 let mut guard = shared.queues.lock().expect("queue mutex poisoned");
1555 if guard.stdout_closed && guard.stderr_closed {
1556 return true;
1557 }
1558 if Instant::now() >= deadline {
1559 guard.stdout_closed = true;
1560 guard.stderr_closed = true;
1561 shared.condvar.notify_all();
1562 return false;
1563 }
1564 }
1565 let remaining = deadline.saturating_duration_since(Instant::now());
1566 tokio::time::sleep(remaining.min(Duration::from_millis(10))).await;
1567 }
1568}
1569
1570fn finalize_capture_completion(shared: &SharedState, deadline: Instant) -> bool {
1571 let mut guard = shared.queues.lock().expect("queue mutex poisoned");
1572 while !(guard.stdout_closed && guard.stderr_closed) {
1573 let now = Instant::now();
1574 if now >= deadline {
1575 guard.stdout_closed = true;
1576 guard.stderr_closed = true;
1577 shared.condvar.notify_all();
1578 return false;
1579 }
1580 let (next_guard, result) = shared
1581 .condvar
1582 .wait_timeout(guard, deadline - now)
1583 .expect("queue mutex poisoned");
1584 guard = next_guard;
1585 if result.timed_out() && !(guard.stdout_closed && guard.stderr_closed) {
1586 guard.stdout_closed = true;
1587 guard.stderr_closed = true;
1588 shared.condvar.notify_all();
1589 return false;
1590 }
1591 }
1592 true
1593}
1594
1595fn emit_lines(shared: &Arc<SharedState>, stream: StreamKind, lines: Vec<Vec<u8>>) {
1596 if lines.is_empty() || shared.capture_overflowed.load(Ordering::Acquire) {
1597 return;
1598 }
1599 let mut guard = shared.queues.lock().expect("queue mutex poisoned");
1600 if shared.capture_overflowed.load(Ordering::Acquire) {
1601 return;
1602 }
1603 for line in lines {
1604 let line_len = line.len();
1605 match stream {
1606 StreamKind::Stdout => {
1607 guard.stdout_history_bytes += line_len;
1608 guard.stdout_history.push_back(line.clone());
1609 guard.stdout_queue.push_back(line.clone());
1610 }
1611 StreamKind::Stderr => {
1612 guard.stderr_history_bytes += line_len;
1613 guard.stderr_history.push_back(line.clone());
1614 guard.stderr_queue.push_back(line.clone());
1615 }
1616 }
1617 let event = StreamEvent { stream, line };
1618 guard.combined_history_bytes += line_len;
1619 guard.combined_history.push_back(event.clone());
1620 guard.combined_queue.push_back(event);
1621 }
1622 shared.condvar.notify_all();
1623}
1624
1625fn append_raw(shared: &Arc<SharedState>, stream: StreamKind, chunk: &[u8]) -> bool {
1626 if chunk.is_empty() {
1627 return true;
1628 }
1629 let mut guard = shared.queues.lock().expect("queue mutex poisoned");
1630 let accepted = match shared.capture_limit {
1631 Some(limit) => {
1632 let retained = guard
1633 .stdout_raw_bytes
1634 .saturating_add(guard.stderr_raw_bytes);
1635 chunk.len().min(limit.saturating_sub(retained))
1636 }
1637 None => chunk.len(),
1638 };
1639 if accepted != 0 {
1640 let accepted_chunk = chunk[..accepted].to_vec();
1641 match stream {
1642 StreamKind::Stdout => {
1643 guard.stdout_raw_bytes += accepted;
1644 guard.stdout_raw.push_back(accepted_chunk);
1645 }
1646 StreamKind::Stderr => {
1647 guard.stderr_raw_bytes += accepted;
1648 guard.stderr_raw.push_back(accepted_chunk);
1649 }
1650 }
1651 }
1652 if accepted != chunk.len() {
1653 shared.capture_overflowed.store(true, Ordering::Release);
1654 false
1655 } else {
1656 shared.condvar.notify_all();
1657 true
1658 }
1659}
1660
1661mod bounded;
1662pub use bounded::{
1663 run_command, run_command_bounded, run_std_command_bounded,
1664 run_std_command_bounded_with_options, BoundedRunOptions,
1665};
1666
1667pub(crate) fn shell_command(command: &str) -> Command {
1668 running_process_platform_internal::platform::process::compat_shell_command(command)
1669}
1670
1671#[cfg(test)]
1672mod tests;