1#![cfg_attr(
10 dylint_lib = "running_process_env_literal",
11 deny(running_process_env_direct)
12)]
13
14use std::cfg_select;
15pub mod env;
17pub mod env_vars;
18pub mod foreground;
19#[cfg(any(windows, test))]
22mod descendant_snapshot;
23mod semantic_priority;
24mod std_child;
25pub use semantic_priority::ProcessPriority;
26pub use std_child::{PlatformCaptureReaders, PlatformStdChild};
27#[cfg(feature = "async-process")]
28mod spawn_admission;
29#[cfg(feature = "async-process")]
30pub use spawn_admission::SpawnAdmission;
31#[cfg(feature = "async-process")]
32use std::ffi::{OsStr, OsString};
33#[cfg(feature = "async-process")]
34use std::io;
35#[cfg(feature = "async-process")]
36use std::path::PathBuf;
37#[cfg(feature = "async-process")]
38use std::process::{ExitStatus, Output, Stdio};
39
40#[cfg(feature = "async-process")]
41use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt};
42#[cfg(feature = "async-process")]
43use tokio::process::{Child, ChildStderr, ChildStdin, ChildStdout, Command};
44
45pub mod platform;
50
51#[cfg(all(feature = "independent-spawn", test, not(windows)))]
52#[path = "platform_win/scheduler_error.rs"]
53mod scheduler_error;
54#[cfg(feature = "independent-spawn")]
55pub(crate) use platform_imp::spawn_sync_owned_daemon;
56#[cfg(feature = "independent-spawn")]
57pub use platform_imp::{
58 independent_broker_run, independent_broker_spawn, independent_spawn, IndependentChild,
59};
60#[cfg(feature = "independent-spawn")]
61pub(crate) use platform_imp::{independent_open_regular, INDEPENDENT_ZERO_WRITE_PENDING};
62
63#[cfg(feature = "pty")]
69#[doc(hidden)]
70pub use portable_pty as portable_pty_compat;
71
72cfg_select! {
75 target_os = "windows" => {
76 mod platform_win;
77 pub(crate) use platform_win as platform_imp;
78 }
79 target_os = "linux" => {
80 mod platform_linux;
81 pub(crate) use platform_linux as platform_imp;
82 }
83 target_os = "macos" => {
84 mod platform_macos;
85 pub(crate) use platform_macos as platform_imp;
86 }
87}
88
89pub(crate) use platform_imp::foreground as foreground_imp;
93pub(crate) use platform_imp::{PRIORITY_NICE_HIGH, PRIORITY_NICE_LOW};
94
95pub use platform_imp::{
96 apply_process_priority, assign_child_to_windows_job, cancel_capture_reader,
97 canonical_environment_pairs, capture_reader_done, compat_shell_command, configure_exact_trace,
98 configure_process_command, configure_process_command_for_bounded_owner_death,
99 configure_sync_contained_command, configure_sync_daemon_command,
100 configure_sync_daemon_command_with_inheritance, configure_trampoline_command,
101 current_executable_build_id, exact_trace_capability, exit_code, exit_signal,
102 monitor_console_windows, parent_has_console, prepare_capture_reader, send_interrupt,
103 set_process_name, shell_command, soft_terminate_process_group, spawn_sync, spawn_sync_daemon,
104 spawn_sync_daemon_with_inheritance, start_attached_descendant_monitor,
105 start_descendant_monitor, start_exact_trace, sync_child_native_handle, trampoline_exit_code,
106 unix_mark_extra_fds_close_on_exec, unix_set_priority, unix_signal_process,
107 unix_signal_process_group, unix_signal_raw, CaptureCancellation, TracedChild, WindowsJobHandle,
108};
109
110#[cfg(feature = "terminal-graphics")]
111pub use platform_imp::active_graphics_probe;
112
113#[cfg(feature = "window-icon")]
114pub use platform_imp::{set_window_icon_impl, window_icon_support_impl};
115
116#[cfg(feature = "async-process")]
117pub(crate) use platform_imp::{
118 async_child_cpu_time, async_child_identity, signal_async_child, signal_async_child_group,
119 AsyncChildIdentity,
120};
121
122#[cfg(feature = "process-inspection")]
123pub use platform_imp::{kill_tree, process_snapshot, process_snapshot_for_pid};
124
125pub use platform_imp::{autostart_register, autostart_render_registration, autostart_unregister};
126
127pub use platform_imp::{process_install_owner_death_cleanup, process_owner_death_cleanup_target};
128
129pub use platform_imp::process_install_shutdown_request_handler;
130
131pub use platform_imp::{fs_open_handles_block_removal, fs_write_all_to_descriptor};
132
133pub use platform_imp::{process_can_replace_current_image, process_replace_current_image};
134
135pub use platform_imp::{process_loaded_images, process_open_loaded_image_file};
136
137#[cfg(feature = "async-process")]
138pub use platform_imp::ape_route_tokio_through_execvp;
139pub use platform_imp::{
140 ape_anonymous_executable, ape_default_loader_dirs, ape_is_exec_format_error, ape_is_executable,
141 ape_mark_executable, ape_private_exec_dir, ape_route_through_execvp, APE_EXECVP_SHELL_FALLBACK,
142 APE_LOADER_HOST, APE_NEEDS_LOADER, APE_SHELL, APE_SYSTEM_LOADERS,
143};
144
145pub use platform_imp::{
146 observer_backend as process_observer_backend, read_process_argv as process_read_argv,
147 read_process_cmdline as process_read_cmdline,
148 read_process_file_handles as process_read_file_handles,
149};
150
151pub use platform_imp::{
152 process_executable_path, process_fault_code_name, process_force_kill,
153 process_same_executable_path, process_signal_terminate, ProcessLiveness,
154};
155
156pub use platform_imp::{
157 resources_fd_exhaustion_error, resources_inode_capacity, resources_signals_fd_exhaustion,
158 resources_signals_storage_exhaustion, resources_storage_exhaustion_error,
159};
160
161pub use platform_imp::{
162 executable_file_name, executable_sibling_of_current_image, EXECUTABLE_EXTENSION,
163};
164
165#[cfg(feature = "fs")]
166pub use platform_imp::{
167 fs_create_private_file, fs_decode_path_bytes, fs_encode_path_bytes, fs_file_identity,
168 fs_is_link_handle, fs_is_lock_conflict, fs_open_lock_file, fs_open_read_no_follow,
169 fs_path_identity, fs_replace_file, fs_state_home_from_environment, fs_sync_directory,
170 fs_try_lock_exclusive, fs_unlock, fs_user_config_dir, fs_user_data_dir, fs_user_run_data_root,
171 fs_user_runtime_dir, fs_user_state_dir, fs_user_state_dir_from_environment, FsFileIdentity,
172};
173
174pub use platform_imp::{
175 host_boot_id, host_current_process_privilege, host_environment_keys_are_case_insensitive,
176 host_filesystem_device_id, host_hostname, host_login_environment, host_machine_id,
177 host_namespace_id, host_process_cgroup, host_user_machine_identity, HostPrivilegedIdentity,
178};
179
180pub use platform_imp::host_login_environment_block;
181
182pub use platform_imp::terminal_input;
183
184#[cfg(feature = "ipc")]
185pub use platform_imp::{
186 ipc_broker_endpoint_name as IpcBrokerEndpointName, ipc_broker_v1_endpoint_path,
187 ipc_broker_v2_runtime_dir, ipc_component_endpoint_path, ipc_component_runtime_dir,
188 ipc_current_user_id, ipc_endpoint_is_filesystem_backed, ipc_endpoint_name_limit,
189 ipc_endpoint_scope_bytes, ipc_handoff_transport_available,
190 ipc_nonblocking_zero_read_is_pending, ipc_select_endpoint_address, IpcEndpoint,
191 IpcInheritedListener, IpcListener, IpcListenerNonblockingMode, IpcPeerIdentity,
192 IpcPeerIdentitySource, IpcStream,
193};
194
195#[cfg(feature = "private-dir")]
196pub use platform_imp::{
197 private_dir_ensure_owner_private_directory, private_dir_owner_private_directory,
198};
199
200#[cfg(feature = "ipc")]
203pub use platform_imp::{
204 private_dir_ensure_owner_private_directory as ipc_ensure_owner_private_directory,
205 private_dir_owner_private_directory as ipc_owner_private_directory,
206};
207
208#[cfg(feature = "ipc")]
213#[doc(hidden)]
214#[derive(Clone, Debug, PartialEq, Eq)]
215pub struct LegacyHandoffError {
216 kind: platform::ipc::HandoffTransferErrorKind,
217 raw_os_error: Option<i32>,
218 transferred_bytes: Option<usize>,
219 expected_bytes: Option<usize>,
220 detail: Option<String>,
221}
222
223#[cfg(feature = "ipc")]
224impl LegacyHandoffError {
225 pub(crate) fn new(
226 kind: platform::ipc::HandoffTransferErrorKind,
227 raw_os_error: Option<i32>,
228 ) -> Self {
229 Self {
230 kind,
231 raw_os_error,
232 transferred_bytes: None,
233 expected_bytes: None,
234 detail: None,
235 }
236 }
237
238 pub(crate) fn with_detail(
239 kind: platform::ipc::HandoffTransferErrorKind,
240 raw_os_error: Option<i32>,
241 detail: impl Into<String>,
242 ) -> Self {
243 Self {
244 kind,
245 raw_os_error,
246 transferred_bytes: None,
247 expected_bytes: None,
248 detail: Some(detail.into()),
249 }
250 }
251
252 #[doc(hidden)]
253 pub fn partial(transferred_bytes: usize, expected_bytes: usize) -> Self {
254 Self {
255 kind: platform::ipc::HandoffTransferErrorKind::Failed,
256 raw_os_error: None,
257 transferred_bytes: Some(transferred_bytes),
258 expected_bytes: Some(expected_bytes),
259 detail: Some(format!(
260 "SCM_RIGHTS connection transfer was partial ({transferred_bytes}/{expected_bytes} bytes)"
261 )),
262 }
263 }
264
265 pub fn kind(&self) -> platform::ipc::HandoffTransferErrorKind {
267 self.kind
268 }
269
270 pub fn raw_os_error(&self) -> Option<i32> {
272 self.raw_os_error
273 }
274
275 pub fn partial_counts(&self) -> Option<(usize, usize)> {
277 self.transferred_bytes.zip(self.expected_bytes)
278 }
279
280 pub(crate) fn detail(&self) -> Option<&str> {
281 self.detail.as_deref()
282 }
283}
284
285#[cfg(feature = "ipc")]
287#[doc(hidden)]
288pub const LEGACY_SCM_RIGHTS_TRANSPORT_SUPPORTED: bool =
289 platform_imp::LEGACY_SCM_RIGHTS_TRANSPORT_SUPPORTED;
290
291#[cfg(feature = "ipc")]
293#[doc(hidden)]
294pub const LEGACY_DUPLICATE_HANDLE_TRANSPORT_SUPPORTED: bool =
295 platform_imp::LEGACY_DUPLICATE_HANDLE_TRANSPORT_SUPPORTED;
296
297#[cfg(feature = "ipc")]
299#[doc(hidden)]
300pub fn legacy_send_fd_to(
301 socket: &std::path::Path,
302 sent_fd: i32,
303 payload: &[u8],
304) -> Result<(), LegacyHandoffError> {
305 platform_imp::legacy_send_fd_to(socket, sent_fd, payload)
306}
307
308#[cfg(feature = "ipc")]
310#[doc(hidden)]
311pub fn legacy_send_fd_over(
312 socket_fd: i32,
313 sent_fd: i32,
314 payload: &[u8],
315) -> Result<(), LegacyHandoffError> {
316 platform_imp::legacy_send_fd_over(socket_fd, sent_fd, payload)
317}
318
319#[cfg(feature = "ipc")]
321#[doc(hidden)]
322pub fn legacy_duplicate_handle(
323 source_handle: usize,
324 backend_pid: u32,
325) -> Result<usize, LegacyHandoffError> {
326 platform_imp::legacy_duplicate_handle(source_handle, backend_pid)
327}
328
329#[cfg(feature = "ipc")]
337#[doc(hidden)]
338pub fn into_legacy_ipc_stream(stream: IpcStream) -> interprocess::local_socket::Stream {
339 platform_imp::into_legacy_ipc_stream(stream)
340}
341
342#[cfg(feature = "ipc")]
344#[doc(hidden)]
345pub fn from_legacy_ipc_stream(stream: interprocess::local_socket::Stream) -> IpcStream {
346 platform_imp::from_legacy_ipc_stream(stream)
347}
348
349#[cfg(feature = "ipc")]
352#[doc(hidden)]
353pub fn legacy_ipc_name(path: &str) -> Result<interprocess::local_socket::Name<'_>, String> {
354 platform_imp::legacy_ipc_name(path)
355}
356
357#[cfg(feature = "ipc-async")]
358pub use platform_imp::{
359 IpcAsyncListener, IpcAsyncStream, IpcIntoAsyncListener, IpcIntoAsyncStream,
360};
361
362#[cfg(feature = "pty")]
363pub use platform_imp::terminal::{
364 before_pty_spawn, current_backend_kind, find_child_processes, find_orphan_conhosts,
365 input_payload, is_ignorable_process_control_error, prepare_unmanaged_pty_child,
366 query_responses, resize_pty, shell_argv, signal_pty_tree, terminate_pty_child,
367 wait_before_pty_close_supported, Backend, ChildProcessInfo, ConPtyBackendKind,
368 OrphanConhostInfo, PtyProcessGuard, PtySpawnContext, TerminalInputSession,
369};
370
371#[cfg(feature = "session-relay")]
372pub use platform_imp::relay_local_socket_session;
373
374#[cfg(feature = "async-process")]
379pub fn configure_compat_tokio_command(
380 command: &mut Command,
381 show_console: bool,
382 kill_when_owner_dies: bool,
383) -> io::Result<()> {
384 platform_imp::configure_compat_tokio_command(command, show_console, kill_when_owner_dies)
385}
386
387#[cfg(feature = "async-process")]
389pub fn after_compat_tokio_spawn(child: &Child, kill_when_owner_dies: bool) -> io::Result<()> {
390 platform_imp::after_compat_tokio_spawn(child, kill_when_owner_dies)
391}
392
393#[cfg(feature = "async-process")]
395#[derive(Debug, Clone, Copy, PartialEq, Eq)]
396pub enum StreamMode {
397 Inherit,
399 Piped,
401 Null,
403}
404
405#[cfg(feature = "async-process")]
406impl StreamMode {
407 fn apply(self) -> Stdio {
408 match self {
409 Self::Inherit => Stdio::inherit(),
410 Self::Piped => Stdio::piped(),
411 Self::Null => Stdio::null(),
412 }
413 }
414}
415
416#[cfg(feature = "async-process")]
418#[derive(Debug, Clone)]
419pub struct SpawnSpec {
420 program: OsString,
421 args: Vec<OsString>,
422 current_dir: Option<PathBuf>,
423 env: Vec<(OsString, OsString)>,
424 clear_env: bool,
425 stdin: StreamMode,
426 stdout: StreamMode,
427 stderr: StreamMode,
428 create_process_group: bool,
429 kill_when_owner_dies: bool,
430 nice: Option<i32>,
431 admission: Option<SpawnAdmission>,
432 command_override: Option<CommandOverride>,
433 creation_flags: Option<u32>,
434 address_space_limit_bytes: Option<u64>,
435}
436
437#[cfg(feature = "async-process")]
442type CommandOverride = std::sync::Arc<std::sync::Mutex<Option<std::process::Command>>>;
443
444#[cfg(feature = "async-process")]
445impl SpawnSpec {
446 pub fn new(program: impl Into<OsString>) -> Self {
448 Self {
449 program: program.into(),
450 args: Vec::new(),
451 current_dir: None,
452 env: Vec::new(),
453 clear_env: false,
454 stdin: StreamMode::Inherit,
455 stdout: StreamMode::Inherit,
456 stderr: StreamMode::Inherit,
457 create_process_group: false,
458 kill_when_owner_dies: false,
459 nice: None,
460 admission: None,
461 command_override: None,
462 creation_flags: None,
463 address_space_limit_bytes: None,
464 }
465 }
466
467 pub fn from_std_command(command: std::process::Command) -> Self {
487 let mut spec = Self::new(command.get_program().to_owned());
488 spec.command_override = Some(std::sync::Arc::new(std::sync::Mutex::new(Some(command))));
489 spec
490 }
491
492 pub fn creation_flags(mut self, flags: Option<u32>) -> Self {
497 self.creation_flags = flags;
498 self
499 }
500
501 pub fn address_space_limit_bytes(mut self, limit: Option<u64>) -> Self {
507 self.address_space_limit_bytes = limit;
508 self
509 }
510
511 pub fn arg(mut self, arg: impl Into<OsString>) -> Self {
513 self.args.push(arg.into());
514 self
515 }
516
517 pub fn current_dir(mut self, path: impl Into<PathBuf>) -> Self {
519 self.current_dir = Some(path.into());
520 self
521 }
522
523 pub fn env(mut self, key: impl Into<OsString>, value: impl Into<OsString>) -> Self {
525 self.env.push((key.into(), value.into()));
526 self
527 }
528
529 pub fn clear_env(mut self, clear: bool) -> Self {
531 self.clear_env = clear;
532 self
533 }
534
535 pub fn stdin(mut self, mode: StreamMode) -> Self {
537 self.stdin = mode;
538 self
539 }
540
541 pub fn stdout(mut self, mode: StreamMode) -> Self {
543 self.stdout = mode;
544 self
545 }
546
547 pub fn stderr(mut self, mode: StreamMode) -> Self {
549 self.stderr = mode;
550 self
551 }
552
553 pub fn create_process_group(mut self, create: bool) -> Self {
562 self.create_process_group = create;
563 self
564 }
565
566 pub fn kill_when_owner_dies(mut self, kill: bool) -> Self {
575 self.kill_when_owner_dies = kill;
576 self
577 }
578
579 pub fn nice(mut self, nice: Option<i32>) -> Self {
586 self.nice = nice;
587 self
588 }
589
590 pub fn priority(self, priority: ProcessPriority) -> Self {
592 self.nice(priority.nice_value())
593 }
594
595 pub fn priority_best_effort(self, priority: ProcessPriority) -> Self {
600 self.priority(priority)
601 }
602
603 pub fn spawn_admission(mut self, admission: SpawnAdmission) -> Self {
605 self.admission = Some(admission);
606 self
607 }
608
609 pub async fn spawn(self) -> io::Result<PlatformChild> {
611 if self.command_override.is_some() {
612 return self.spawn_override();
613 }
614 if self.creation_flags.is_some() || self.address_space_limit_bytes.is_some() {
615 return Err(io::Error::new(
616 io::ErrorKind::InvalidInput,
617 "creation flags and address-space limits apply only to SpawnSpec::from_std_command",
618 ));
619 }
620 let options = platform::ape::ApeOptions::with_overrides(
626 self.clear_env,
627 self.env
628 .iter()
629 .map(|(key, value)| (key.as_os_str(), Some(value.as_os_str()))),
630 );
631 let mut command = match platform::ape::plan_launch(
632 &self.program,
633 self.current_dir.as_deref(),
634 &options,
635 ) {
636 Some(launch) => {
637 let mut command =
638 self.command(launch.loader.as_os_str(), &launch.args(&self.args))?;
639 if let Some(path) = launch.child_path(options.path.as_deref()) {
640 command.env("PATH", path);
641 }
642 command
643 }
644 None => self.command(&self.program, &self.args)?,
645 };
646 let mut spawn = || {
649 platform::ape::retry_while_busy(|| {
650 let _fork = platform::ape::fork_guard();
651 command.spawn()
652 })
653 };
654 let child = match self.admission.as_ref() {
655 Some(admission) => admission.run(spawn)?,
656 None => spawn()?,
657 };
658 platform_imp::after_spawn(&child, self.kill_when_owner_dies, self.nice)?;
659 Ok(PlatformChild::new(child, self.create_process_group))
660 }
661
662 fn command(&self, program: &OsStr, args: &[OsString]) -> io::Result<Command> {
664 let mut command = Command::new(program);
665 command.args(args);
666 if let Some(current_dir) = self.current_dir.as_deref() {
667 command.current_dir(current_dir);
668 }
669 if self.clear_env {
670 command.env_clear();
671 }
672 for (key, value) in &self.env {
673 command.env(key, value);
674 }
675 command
676 .stdin(self.stdin.apply())
677 .stdout(self.stdout.apply())
678 .stderr(self.stderr.apply());
679 platform_imp::configure_command(
680 &mut command,
681 self.create_process_group,
682 self.kill_when_owner_dies,
683 self.nice,
684 )?;
685 Ok(command)
686 }
687
688 fn spawn_override(self) -> io::Result<PlatformChild> {
689 if !self.args.is_empty()
690 || !self.env.is_empty()
691 || self.clear_env
692 || self.current_dir.is_some()
693 {
694 return Err(io::Error::new(
695 io::ErrorKind::InvalidInput,
696 "arguments, environment and working directory belong on the caller-built command",
697 ));
698 }
699 let mut std_command = self
700 .command_override
701 .as_ref()
702 .expect("override checked by caller")
703 .lock()
704 .unwrap_or_else(std::sync::PoisonError::into_inner)
705 .take()
706 .ok_or_else(|| {
707 io::Error::new(
708 io::ErrorKind::InvalidInput,
709 "caller-built command was already spawned by a clone of this spec",
710 )
711 })?;
712 platform_imp::configure_override_command(
713 &mut std_command,
714 platform::process::ProcessCommandConfig {
715 creation_flags: self.creation_flags,
716 create_process_group: self.create_process_group,
717 nice: self.nice,
718 address_space_limit_bytes: self.address_space_limit_bytes,
719 },
720 self.kill_when_owner_dies,
721 )?;
722 let mut command = Command::from(std_command);
723 command
724 .stdin(self.stdin.apply())
725 .stdout(self.stdout.apply())
726 .stderr(self.stderr.apply());
727 let mut spawn = || platform::ape::spawn_tokio(&mut command, |command| command.spawn());
728 let child = match self.admission.as_ref() {
729 Some(admission) => admission.run(spawn)?,
730 None => spawn()?,
731 };
732 Ok(PlatformChild::new(child, self.create_process_group))
733 }
734}
735
736#[cfg(feature = "async-process")]
738pub struct PlatformChild {
739 child: Child,
740 stdin: Option<ChildStdin>,
741 stdout: Option<ChildStdout>,
742 stderr: Option<ChildStderr>,
743 signal: PlatformEmergencySignal,
744}
745
746#[cfg(feature = "async-process")]
747impl PlatformChild {
748 fn new(mut child: Child, own_process_group: bool) -> Self {
749 let signal = PlatformEmergencySignal {
750 identity: async_child_identity(&child),
751 own_process_group,
752 legacy_group_pid: child.id(),
757 };
758 Self {
759 stdin: child.stdin.take(),
760 stdout: child.stdout.take(),
761 stderr: child.stderr.take(),
762 child,
763 signal,
764 }
765 }
766
767 pub fn id(&self) -> Option<u32> {
769 self.child.id()
770 }
771
772 pub async fn wait(&mut self) -> io::Result<ExitStatus> {
774 self.child.wait().await
775 }
776
777 pub async fn kill(&mut self) -> io::Result<()> {
779 self.child.kill().await
780 }
781
782 pub async fn wait_with_output(self) -> io::Result<Output> {
784 let Self {
785 mut child,
786 stdin,
787 stdout,
788 stderr,
789 ..
790 } = self;
791 drop(stdin);
794 let (status, stdout, stderr) = tokio::try_join!(
795 child.wait(),
796 read_owned_to_end(stdout),
797 read_owned_to_end(stderr),
798 )?;
799 Ok(Output {
800 status,
801 stdout,
802 stderr,
803 })
804 }
805
806 pub async fn write_stdin(&mut self, bytes: &[u8]) -> io::Result<()> {
808 let stdin = self.stdin.as_mut().ok_or_else(stdin_not_piped)?;
809 stdin.write_all(bytes).await?;
810 stdin.flush().await
811 }
812
813 pub fn close_stdin(&mut self) {
818 drop(self.stdin.take());
819 }
820
821 pub async fn read_stdout_to_end(&mut self) -> io::Result<Vec<u8>> {
823 let stdout = self.stdout.as_mut().ok_or_else(stdout_not_piped)?;
824 let mut bytes = Vec::new();
825 stdout.read_to_end(&mut bytes).await?;
826 Ok(bytes)
827 }
828
829 pub async fn read_stderr_to_end(&mut self) -> io::Result<Vec<u8>> {
831 let stderr = self.stderr.as_mut().ok_or_else(stderr_not_piped)?;
832 let mut bytes = Vec::new();
833 stderr.read_to_end(&mut bytes).await?;
834 Ok(bytes)
835 }
836
837 pub fn into_actor_parts(
843 self,
844 ) -> (
845 PlatformLifecycle,
846 PlatformEmergencySignal,
847 Option<PlatformStdin>,
848 Option<PlatformOutput>,
849 Option<PlatformOutput>,
850 ) {
851 (
852 PlatformLifecycle { child: self.child },
853 self.signal,
854 self.stdin.map(|stdin| PlatformStdin { stdin }),
855 self.stdout.map(PlatformOutput::stdout),
856 self.stderr.map(PlatformOutput::stderr),
857 )
858 }
859}
860
861#[cfg(feature = "async-process")]
863pub struct PlatformLifecycle {
864 child: Child,
865}
866
867#[cfg(feature = "async-process")]
868impl PlatformLifecycle {
869 pub async fn wait(&mut self) -> io::Result<ExitStatus> {
871 self.child.wait().await
872 }
873
874 pub fn start_kill(&mut self) -> io::Result<()> {
882 self.child.start_kill()
883 }
884}
885
886#[cfg(feature = "async-process")]
891pub struct PlatformEmergencySignal {
892 identity: Option<AsyncChildIdentity>,
893 own_process_group: bool,
894 legacy_group_pid: Option<u32>,
895}
896
897#[cfg(feature = "async-process")]
898impl PlatformEmergencySignal {
899 pub fn kill(&self) -> io::Result<()> {
901 let Some(identity) = self.identity.as_ref() else {
902 return Err(signal_target_unavailable());
903 };
904 signal_async_child(identity)
905 }
906
907 pub fn terminate_group_soft(&self) -> io::Result<bool> {
916 if !self.own_process_group {
917 return Ok(false);
918 }
919 let Some(identity) = self.identity.as_ref() else {
920 return Err(signal_target_unavailable());
921 };
922 signal_async_child_group(identity).map(|()| true)
923 }
924
925 pub fn terminate_group_soft_legacy(&self) -> io::Result<bool> {
934 if !self.own_process_group {
935 return Ok(false);
936 }
937 if let Some(identity) = self.identity.as_ref() {
938 return signal_async_child_group(identity).map(|()| true);
939 }
940 let pid = self
941 .legacy_group_pid
942 .ok_or_else(signal_target_unavailable)?;
943 crate::platform::process::soft_terminate_process_group(pid).map(|()| true)
944 }
945
946 pub fn cpu_time(&self) -> io::Result<Option<std::time::Duration>> {
950 self.identity
951 .as_ref()
952 .map_or(Ok(None), async_child_cpu_time)
953 }
954}
955
956#[cfg(feature = "async-process")]
957fn signal_target_unavailable() -> io::Error {
958 io::Error::new(
959 io::ErrorKind::BrokenPipe,
960 "child process launch identity is no longer available",
961 )
962}
963
964#[cfg(feature = "async-process")]
966pub struct PlatformStdin {
967 stdin: ChildStdin,
968}
969
970#[cfg(feature = "async-process")]
971impl PlatformStdin {
972 pub async fn write(&mut self, bytes: &[u8]) -> io::Result<()> {
974 self.stdin.write_all(bytes).await?;
975 self.stdin.flush().await
976 }
977}
978
979#[cfg(feature = "async-process")]
981pub struct PlatformOutput {
982 reader: OutputReader,
983 read_pending: bool,
984}
985
986#[cfg(feature = "async-process")]
987enum OutputReader {
988 Stdout(ChildStdout),
989 Stderr(ChildStderr),
990}
991
992#[cfg(feature = "async-process")]
993impl PlatformOutput {
994 fn stdout(stdout: ChildStdout) -> Self {
995 Self {
996 reader: OutputReader::Stdout(stdout),
997 read_pending: false,
998 }
999 }
1000
1001 fn stderr(stderr: ChildStderr) -> Self {
1002 Self {
1003 reader: OutputReader::Stderr(stderr),
1004 read_pending: false,
1005 }
1006 }
1007
1008 pub async fn read_to_end(self) -> io::Result<Vec<u8>> {
1010 match self.reader {
1011 OutputReader::Stdout(stdout) => read_owned_to_end(Some(stdout)).await,
1012 OutputReader::Stderr(stderr) => read_owned_to_end(Some(stderr)).await,
1013 }
1014 }
1015
1016 pub async fn read_chunk(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
1021 self.read_pending = true;
1022 let result = match &mut self.reader {
1023 OutputReader::Stdout(stdout) => stdout.read(buffer).await,
1024 OutputReader::Stderr(stderr) => stderr.read(buffer).await,
1025 };
1026 self.read_pending = false;
1027 result
1028 }
1029
1030 pub async fn shutdown(self) -> io::Result<()> {
1037 match self.reader {
1038 OutputReader::Stdout(stdout) => {
1039 platform_imp::shutdown_output_reader(stdout, self.read_pending).await
1040 }
1041 OutputReader::Stderr(stderr) => {
1042 platform_imp::shutdown_output_reader(stderr, self.read_pending).await
1043 }
1044 }
1045 }
1046}
1047
1048#[cfg(all(test, feature = "async-process"))]
1049mod output_shutdown_tests {
1050 use super::*;
1051 use std::time::Duration;
1052
1053 #[test]
1054 fn silent_output_fixture() {
1055 if std::env::var_os("RUNNING_PROCESS_OUTPUT_SHUTDOWN_FIXTURE").is_some() {
1056 std::thread::sleep(Duration::from_secs(30));
1057 }
1058 }
1059
1060 #[tokio::test]
1061 async fn shutdown_finishes_a_cancelled_pending_read_before_child_exit() {
1062 let child = SpawnSpec::new(std::env::current_exe().expect("test executable"))
1063 .arg("--exact")
1064 .arg("output_shutdown_tests::silent_output_fixture")
1065 .env("RUNNING_PROCESS_OUTPUT_SHUTDOWN_FIXTURE", "1")
1066 .stdin(StreamMode::Null)
1067 .stdout(StreamMode::Piped)
1068 .stderr(StreamMode::Null)
1069 .spawn()
1070 .await
1071 .expect("spawn silent fixture");
1072 let (mut lifecycle, _, _, stdout, _) = child.into_actor_parts();
1073 let mut stdout = stdout.expect("stdout pipe");
1074 let mut bytes = [0; 1024];
1076 loop {
1077 match tokio::time::timeout(Duration::from_millis(20), stdout.read_chunk(&mut bytes))
1078 .await
1079 {
1080 Ok(Ok(0)) => panic!("fixture exited before pending read"),
1081 Ok(Ok(_)) => {}
1082 Ok(Err(error)) => panic!("fixture read failed: {error}"),
1083 Err(_) => break,
1084 }
1085 }
1086 let shutdown = tokio::time::timeout(Duration::from_secs(2), stdout.shutdown()).await;
1087 lifecycle.start_kill().expect("kill fixture");
1090 lifecycle.wait().await.expect("reap fixture");
1091 shutdown
1092 .expect("pending read shutdown must finish")
1093 .expect("output shutdown");
1094 }
1095}
1096
1097#[cfg(feature = "async-process")]
1098fn stdin_not_piped() -> io::Error {
1099 io::Error::new(io::ErrorKind::BrokenPipe, "child stdin is not piped")
1100}
1101
1102#[cfg(feature = "async-process")]
1103fn stdout_not_piped() -> io::Error {
1104 io::Error::new(io::ErrorKind::BrokenPipe, "child stdout is not piped")
1105}
1106
1107#[cfg(feature = "async-process")]
1108fn stderr_not_piped() -> io::Error {
1109 io::Error::new(io::ErrorKind::BrokenPipe, "child stderr is not piped")
1110}
1111
1112#[cfg(feature = "async-process")]
1113async fn read_owned_to_end<R>(reader: Option<R>) -> io::Result<Vec<u8>>
1114where
1115 R: AsyncRead + Unpin,
1116{
1117 let Some(mut reader) = reader else {
1118 return Ok(Vec::new());
1119 };
1120 let mut bytes = Vec::new();
1121 reader.read_to_end(&mut bytes).await?;
1122 Ok(bytes)
1123}
1124
1125#[cfg(feature = "async-process")]
1127pub fn shell_spec(command: impl AsRef<OsStr>) -> SpawnSpec {
1128 platform_imp::shell_spec(command.as_ref())
1129}
1130
1131#[cfg(all(test, feature = "async-process"))]
1132mod tests {
1133 use super::{shell_spec, SpawnSpec, StreamMode};
1134
1135 fn fixture_command() -> SpawnSpec {
1136 #[cfg(windows)]
1137 {
1138 shell_spec("echo async-platform-internal")
1139 }
1140 #[cfg(not(windows))]
1141 {
1142 shell_spec("printf async-platform-internal")
1143 }
1144 }
1145
1146 #[tokio::test]
1147 async fn blessed_spawn_captures_output_without_sync_wait() {
1148 let output = fixture_command()
1149 .stdout(StreamMode::Piped)
1150 .stderr(StreamMode::Piped)
1151 .spawn()
1152 .await
1153 .expect("spawn")
1154 .wait_with_output()
1155 .await
1156 .expect("wait with output");
1157
1158 assert!(output.status.success());
1159 let expected = if cfg!(windows) {
1160 b"async-platform-internal\r\n".as_slice()
1161 } else {
1162 b"async-platform-internal".as_slice()
1163 };
1164 assert_eq!(output.stdout, expected);
1165 assert!(output.stderr.is_empty());
1166 }
1167
1168 #[test]
1170 fn override_env_fixture() {
1171 if std::env::var_os("RUNNING_PROCESS_OVERRIDE_FIXTURE").is_none() {
1172 return;
1173 }
1174 let path = if std::env::var_os("PATH").is_some() {
1175 "present"
1176 } else {
1177 "absent"
1178 };
1179 println!("override-fixture PATH={path}");
1180 }
1181
1182 fn override_fixture_command() -> std::process::Command {
1183 let mut command =
1184 std::process::Command::new(std::env::current_exe().expect("test executable"));
1185 command
1186 .args(["--exact", "tests::override_env_fixture", "--nocapture"])
1187 .env("RUNNING_PROCESS_OVERRIDE_FIXTURE", "1");
1188 command
1189 }
1190
1191 async fn override_fixture_stdout(command: std::process::Command) -> String {
1192 let output = SpawnSpec::from_std_command(command)
1193 .stdin(StreamMode::Null)
1194 .stdout(StreamMode::Piped)
1195 .stderr(StreamMode::Null)
1196 .spawn()
1197 .await
1198 .expect("spawn override")
1199 .wait_with_output()
1200 .await
1201 .expect("wait override");
1202 assert!(output.status.success());
1203 String::from_utf8_lossy(&output.stdout).into_owned()
1204 }
1205
1206 #[tokio::test]
1207 async fn command_override_keeps_inherited_environment_by_default() {
1208 let stdout = override_fixture_stdout(override_fixture_command()).await;
1209 assert!(stdout.contains("override-fixture PATH=present"), "{stdout}");
1210 }
1211
1212 #[tokio::test]
1213 async fn command_override_preserves_env_remove_of_inherited_variable() {
1214 let mut command = override_fixture_command();
1215 command.env_remove("PATH");
1216 let stdout = override_fixture_stdout(command).await;
1217 assert!(stdout.contains("override-fixture PATH=absent"), "{stdout}");
1218 }
1219
1220 #[tokio::test]
1221 async fn command_override_rejects_declarative_command_fields() {
1222 for spec in [
1223 SpawnSpec::from_std_command(override_fixture_command()).arg("extra"),
1224 SpawnSpec::from_std_command(override_fixture_command()).env("K", "V"),
1225 SpawnSpec::from_std_command(override_fixture_command()).clear_env(true),
1226 SpawnSpec::from_std_command(override_fixture_command()).current_dir("."),
1227 ] {
1228 let error = spec.spawn().await.err().expect("must be rejected");
1229 assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput);
1230 }
1231 }
1232
1233 #[tokio::test]
1234 async fn command_override_is_spawned_by_exactly_one_clone() {
1235 let spec = SpawnSpec::from_std_command(override_fixture_command())
1236 .stdin(StreamMode::Null)
1237 .stdout(StreamMode::Null)
1238 .stderr(StreamMode::Null);
1239 let clone = spec.clone();
1240 let mut child = spec.spawn().await.expect("first spawn");
1241 let error = clone.spawn().await.err().expect("second spawn refused");
1242 assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput);
1243 assert!(child.wait().await.expect("reap").success());
1244 }
1245
1246 #[tokio::test]
1247 async fn declarative_spawn_rejects_override_only_options() {
1248 for spec in [
1249 fixture_command().creation_flags(Some(0)),
1250 fixture_command().address_space_limit_bytes(Some(1 << 40)),
1251 ] {
1252 let error = spec.spawn().await.err().expect("must be rejected");
1253 assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput);
1254 }
1255 }
1256
1257 #[tokio::test]
1258 async fn blessed_spawn_reports_missing_program() {
1259 let result = SpawnSpec::new("running-process-program-that-does-not-exist")
1260 .spawn()
1261 .await;
1262 assert!(result.is_err());
1263 }
1264
1265 #[tokio::test]
1266 async fn one_shot_output_closes_owned_stdin() {
1267 #[cfg(windows)]
1268 let spec = shell_spec("more > nul & echo done");
1269 #[cfg(not(windows))]
1270 let spec = shell_spec("cat > /dev/null; printf done");
1271
1272 let output = tokio::time::timeout(
1273 std::time::Duration::from_secs(2),
1274 spec.stdin(StreamMode::Piped)
1275 .stdout(StreamMode::Piped)
1276 .stderr(StreamMode::Piped)
1277 .spawn()
1278 .await
1279 .expect("spawn")
1280 .wait_with_output(),
1281 )
1282 .await
1283 .expect("stdin is closed for one-shot output")
1284 .expect("output succeeds");
1285
1286 let expected = if cfg!(windows) {
1287 b"done\r\n".as_slice()
1288 } else {
1289 b"done".as_slice()
1290 };
1291 assert_eq!(output.stdout, expected);
1292 }
1293}