Skip to main content

running_process_platform_internal/
lib.rs

1//! Blessed asynchronous process operations.
2//!
3//! This crate is intentionally published as an implementation detail. It is
4//! the only production owner of the Tokio process primitives used by the
5//! async process API. Higher layers receive typed operations and never name
6//! `tokio::process::Command` directly.
7
8use std::cfg_select;
9/// Explicit caller-owned foreground command execution.
10pub mod foreground;
11mod semantic_priority;
12pub use semantic_priority::ProcessPriority;
13#[cfg(feature = "async-process")]
14mod spawn_admission;
15#[cfg(feature = "async-process")]
16pub use spawn_admission::SpawnAdmission;
17#[cfg(feature = "async-process")]
18use std::ffi::{OsStr, OsString};
19#[cfg(feature = "async-process")]
20use std::io;
21#[cfg(feature = "async-process")]
22use std::path::PathBuf;
23#[cfg(feature = "async-process")]
24use std::process::{ExitStatus, Output, Stdio};
25
26#[cfg(feature = "async-process")]
27use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt};
28#[cfg(feature = "async-process")]
29use tokio::process::{Child, ChildStderr, ChildStdin, ChildStdout, Command};
30
31/// Neutral capability indexes for the eventual workspace-wide host boundary.
32///
33/// The indexes intentionally expose no operations yet: phase 2 establishes
34/// ownership names before later phases move a capability behind them.
35pub mod platform;
36
37#[cfg(all(feature = "independent-spawn", test, not(windows)))]
38#[path = "platform_win/scheduler_error.rs"]
39mod scheduler_error;
40#[cfg(feature = "independent-spawn")]
41pub(crate) use platform_imp::spawn_sync_owned_daemon;
42#[cfg(feature = "independent-spawn")]
43pub use platform_imp::{
44    independent_broker_run, independent_broker_spawn, independent_spawn, IndependentChild,
45};
46#[cfg(feature = "independent-spawn")]
47pub(crate) use platform_imp::{independent_open_regular, INDEPENDENT_ZERO_WRITE_PENDING};
48
49/// Temporary source-compatibility re-export for the pre-boundary PTY API.
50///
51/// New code must use [`platform::terminal`] facade-owned types. This root-only
52/// alias deliberately stays outside the neutral facade and can be removed in
53/// the next major release after downstream users have migrated.
54#[cfg(feature = "pty")]
55#[doc(hidden)]
56pub use portable_pty as portable_pty_compat;
57
58// This is deliberately the crate's only host selector.  Facade modules are
59// neutral; native details live behind the selected private root.
60cfg_select! {
61    target_os = "windows" => {
62        mod platform_win;
63        pub(crate) use platform_win as platform_imp;
64    }
65    target_os = "linux" => {
66        mod platform_linux;
67        pub(crate) use platform_linux as platform_imp;
68    }
69    target_os = "macos" => {
70        mod platform_macos;
71        pub(crate) use platform_macos as platform_imp;
72    }
73}
74
75// Re-export the selected implementation once from this allowed host-selector
76// root. Neutral capability facades re-export only crate-root names and never
77// name the private `platform_imp` alias themselves.
78pub(crate) use platform_imp::foreground as foreground_imp;
79pub(crate) use platform_imp::{PRIORITY_NICE_HIGH, PRIORITY_NICE_LOW};
80
81pub use platform_imp::{
82    assign_child_to_windows_job, cancel_capture_reader, canonical_environment_pairs,
83    capture_reader_done, compat_shell_command, configure_exact_trace, configure_process_command,
84    configure_process_command_for_bounded_owner_death, configure_sync_contained_command,
85    configure_sync_daemon_command, configure_sync_daemon_command_with_inheritance,
86    configure_trampoline_command, current_executable_build_id, exact_trace_capability, exit_code,
87    monitor_console_windows, parent_has_console, prepare_capture_reader, set_process_name,
88    shell_command, soft_terminate_process_group, spawn_sync, spawn_sync_daemon,
89    spawn_sync_daemon_with_inheritance, start_descendant_monitor, start_exact_trace,
90    sync_child_native_handle, trampoline_exit_code, unix_mark_extra_fds_close_on_exec,
91    unix_set_priority, unix_signal_process, unix_signal_process_group, unix_signal_raw,
92    CaptureCancellation, TracedChild, WindowsJobHandle,
93};
94
95#[cfg(feature = "terminal-graphics")]
96pub use platform_imp::active_graphics_probe;
97
98#[cfg(feature = "window-icon")]
99pub use platform_imp::{set_window_icon_impl, window_icon_support_impl};
100
101#[cfg(feature = "async-process")]
102pub(crate) use platform_imp::{
103    async_child_cpu_time, async_child_identity, signal_async_child, signal_async_child_group,
104    AsyncChildIdentity,
105};
106
107#[cfg(feature = "process-inspection")]
108pub use platform_imp::{kill_tree, process_snapshot, process_snapshot_for_pid};
109
110pub use platform_imp::{autostart_register, autostart_render_registration, autostart_unregister};
111
112pub use platform_imp::{process_install_owner_death_cleanup, process_owner_death_cleanup_target};
113
114pub use platform_imp::process_install_shutdown_request_handler;
115
116pub use platform_imp::fs_write_all_to_descriptor;
117
118pub use platform_imp::{process_can_replace_current_image, process_replace_current_image};
119
120pub use platform_imp::{
121    process_executable_path, process_force_kill, process_same_executable_path,
122    process_signal_terminate, ProcessLiveness,
123};
124
125pub use platform_imp::{
126    resources_fd_exhaustion_error, resources_inode_capacity, resources_signals_fd_exhaustion,
127    resources_signals_storage_exhaustion, resources_storage_exhaustion_error,
128};
129
130pub use platform_imp::{
131    executable_file_name, executable_sibling_of_current_image, EXECUTABLE_EXTENSION,
132};
133
134#[cfg(feature = "fs")]
135pub use platform_imp::{
136    fs_create_private_file, fs_decode_path_bytes, fs_encode_path_bytes, fs_file_identity,
137    fs_is_lock_conflict, fs_open_lock_file, fs_path_identity, fs_replace_file, fs_sync_directory,
138    fs_try_lock_exclusive, fs_unlock, fs_user_config_dir, fs_user_data_dir, fs_user_run_data_root,
139    fs_user_runtime_dir, fs_user_state_dir, FsFileIdentity,
140};
141
142pub use platform_imp::{
143    host_boot_id, host_current_process_privilege, host_environment_keys_are_case_insensitive,
144    host_filesystem_device_id, host_hostname, host_login_environment, host_machine_id,
145    host_namespace_id, host_user_machine_identity, HostPrivilegedIdentity,
146};
147
148pub use platform_imp::host_login_environment_block;
149
150pub use platform_imp::terminal_input;
151
152#[cfg(feature = "ipc")]
153pub use platform_imp::{
154    ipc_broker_endpoint_name as IpcBrokerEndpointName, ipc_broker_v1_endpoint_path,
155    ipc_broker_v2_runtime_dir, ipc_current_user_id, ipc_endpoint_is_filesystem_backed,
156    ipc_endpoint_name_limit, ipc_endpoint_scope_bytes, ipc_nonblocking_zero_read_is_pending,
157    ipc_select_endpoint_address, IpcEndpoint, IpcInheritedListener, IpcListener,
158    IpcListenerNonblockingMode, IpcPeerIdentity, IpcPeerIdentitySource, IpcStream,
159};
160
161#[cfg(feature = "private-dir")]
162pub use platform_imp::{
163    private_dir_ensure_owner_private_directory, private_dir_owner_private_directory,
164};
165
166// Retain the implementation-detail root aliases selected by the historical
167// `ipc` capability.  New callers use `platform::private_dir` instead.
168#[cfg(feature = "ipc")]
169pub use platform_imp::{
170    private_dir_ensure_owner_private_directory as ipc_ensure_owner_private_directory,
171    private_dir_owner_private_directory as ipc_owner_private_directory,
172};
173
174/// Failure details for the deprecated 4.x raw descriptor/handle handoff API.
175///
176/// This type exists only at the crate-root compatibility boundary. New product
177/// mechanics use opaque [`platform::ipc::Stream`] operations instead.
178#[cfg(feature = "ipc")]
179#[doc(hidden)]
180#[derive(Clone, Debug, PartialEq, Eq)]
181pub struct LegacyHandoffError {
182    kind: platform::ipc::HandoffTransferErrorKind,
183    raw_os_error: Option<i32>,
184    transferred_bytes: Option<usize>,
185    expected_bytes: Option<usize>,
186    detail: Option<String>,
187}
188
189#[cfg(feature = "ipc")]
190impl LegacyHandoffError {
191    pub(crate) fn new(
192        kind: platform::ipc::HandoffTransferErrorKind,
193        raw_os_error: Option<i32>,
194    ) -> Self {
195        Self {
196            kind,
197            raw_os_error,
198            transferred_bytes: None,
199            expected_bytes: None,
200            detail: None,
201        }
202    }
203
204    pub(crate) fn with_detail(
205        kind: platform::ipc::HandoffTransferErrorKind,
206        raw_os_error: Option<i32>,
207        detail: impl Into<String>,
208    ) -> Self {
209        Self {
210            kind,
211            raw_os_error,
212            transferred_bytes: None,
213            expected_bytes: None,
214            detail: Some(detail.into()),
215        }
216    }
217
218    #[doc(hidden)]
219    pub fn partial(transferred_bytes: usize, expected_bytes: usize) -> Self {
220        Self {
221            kind: platform::ipc::HandoffTransferErrorKind::Failed,
222            raw_os_error: None,
223            transferred_bytes: Some(transferred_bytes),
224            expected_bytes: Some(expected_bytes),
225            detail: Some(format!(
226                "SCM_RIGHTS connection transfer was partial ({transferred_bytes}/{expected_bytes} bytes)"
227            )),
228        }
229    }
230
231    /// Return the policy-neutral failure category.
232    pub fn kind(&self) -> platform::ipc::HandoffTransferErrorKind {
233        self.kind
234    }
235
236    /// Return the native error code retained for legacy public diagnostics.
237    pub fn raw_os_error(&self) -> Option<i32> {
238        self.raw_os_error
239    }
240
241    /// Return a partial payload count when the descriptor may have transferred.
242    pub fn partial_counts(&self) -> Option<(usize, usize)> {
243        self.transferred_bytes.zip(self.expected_bytes)
244    }
245
246    pub(crate) fn detail(&self) -> Option<&str> {
247        self.detail.as_deref()
248    }
249}
250
251/// Whether the deprecated 4.x SCM_RIGHTS compatibility transport is available.
252#[cfg(feature = "ipc")]
253#[doc(hidden)]
254pub const LEGACY_SCM_RIGHTS_TRANSPORT_SUPPORTED: bool =
255    platform_imp::LEGACY_SCM_RIGHTS_TRANSPORT_SUPPORTED;
256
257/// Whether the deprecated 4.x DuplicateHandle compatibility transport is available.
258#[cfg(feature = "ipc")]
259#[doc(hidden)]
260pub const LEGACY_DUPLICATE_HANDLE_TRANSPORT_SUPPORTED: bool =
261    platform_imp::LEGACY_DUPLICATE_HANDLE_TRANSPORT_SUPPORTED;
262
263/// Root-only adapter for the deprecated raw-descriptor handoff API.
264#[cfg(feature = "ipc")]
265#[doc(hidden)]
266pub fn legacy_send_fd_to(
267    socket: &std::path::Path,
268    sent_fd: i32,
269    payload: &[u8],
270) -> Result<(), LegacyHandoffError> {
271    platform_imp::legacy_send_fd_to(socket, sent_fd, payload)
272}
273
274/// Root-only adapter for the deprecated connected raw-descriptor handoff API.
275#[cfg(feature = "ipc")]
276#[doc(hidden)]
277pub fn legacy_send_fd_over(
278    socket_fd: i32,
279    sent_fd: i32,
280    payload: &[u8],
281) -> Result<(), LegacyHandoffError> {
282    platform_imp::legacy_send_fd_over(socket_fd, sent_fd, payload)
283}
284
285/// Root-only adapter for the deprecated raw-handle duplication API.
286#[cfg(feature = "ipc")]
287#[doc(hidden)]
288pub fn legacy_duplicate_handle(
289    source_handle: usize,
290    backend_pid: u32,
291) -> Result<usize, LegacyHandoffError> {
292    platform_imp::legacy_duplicate_handle(source_handle, backend_pid)
293}
294
295/// Temporary source-compatibility conversion for public APIs that predate the
296/// opaque IPC facade.
297///
298/// New code must keep [`IpcStream`] opaque. This root-only adapter exists so
299/// `running-process` can preserve its established raw-stream callback contract
300/// until the next major release without exposing the transport through
301/// [`platform::ipc`].
302#[cfg(feature = "ipc")]
303#[doc(hidden)]
304pub fn into_legacy_ipc_stream(stream: IpcStream) -> interprocess::local_socket::Stream {
305    platform_imp::into_legacy_ipc_stream(stream)
306}
307
308/// Temporary source-compatibility conversion for legacy 4.x callback inputs.
309#[cfg(feature = "ipc")]
310#[doc(hidden)]
311pub fn from_legacy_ipc_stream(stream: interprocess::local_socket::Stream) -> IpcStream {
312    platform_imp::from_legacy_ipc_stream(stream)
313}
314
315/// Temporary source-compatibility conversion for public APIs that return an
316/// `interprocess` endpoint name.
317#[cfg(feature = "ipc")]
318#[doc(hidden)]
319pub fn legacy_ipc_name(path: &str) -> Result<interprocess::local_socket::Name<'_>, String> {
320    platform_imp::legacy_ipc_name(path)
321}
322
323#[cfg(feature = "ipc-async")]
324pub use platform_imp::{
325    IpcAsyncListener, IpcAsyncStream, IpcIntoAsyncListener, IpcIntoAsyncStream,
326};
327
328#[cfg(feature = "pty")]
329pub use platform_imp::terminal::{
330    before_pty_spawn, current_backend_kind, find_child_processes, find_orphan_conhosts,
331    input_payload, is_ignorable_process_control_error, prepare_unmanaged_pty_child,
332    query_responses, resize_pty, shell_argv, signal_pty_tree, terminate_pty_child,
333    wait_before_pty_close_supported, Backend, ChildProcessInfo, ConPtyBackendKind,
334    OrphanConhostInfo, PtyProcessGuard, PtySpawnContext, TerminalInputSession,
335};
336
337#[cfg(feature = "session-relay")]
338pub use platform_imp::relay_local_socket_session;
339
340/// Apply host-owned setup for the legacy Tokio-command compatibility surface.
341///
342/// The public wrapper retains its policy type, while console suppression and
343/// owner-death primitives stay inside the selected platform root.
344#[cfg(feature = "async-process")]
345pub fn configure_compat_tokio_command(
346    command: &mut Command,
347    show_console: bool,
348    kill_when_owner_dies: bool,
349) -> io::Result<()> {
350    platform_imp::configure_compat_tokio_command(command, show_console, kill_when_owner_dies)
351}
352
353/// Complete host-owned setup after a legacy Tokio child has been spawned.
354#[cfg(feature = "async-process")]
355pub fn after_compat_tokio_spawn(child: &Child, kill_when_owner_dies: bool) -> io::Result<()> {
356    platform_imp::after_compat_tokio_spawn(child, kill_when_owner_dies)
357}
358
359/// Stdio policy for one child stream.
360#[cfg(feature = "async-process")]
361#[derive(Debug, Clone, Copy, PartialEq, Eq)]
362pub enum StreamMode {
363    /// Leave the stream connected to the parent process.
364    Inherit,
365    /// Create an asynchronous pipe owned by the child handle.
366    Piped,
367    /// Connect the stream to the platform null device.
368    Null,
369}
370
371#[cfg(feature = "async-process")]
372impl StreamMode {
373    fn apply(self) -> Stdio {
374        match self {
375            Self::Inherit => Stdio::inherit(),
376            Self::Piped => Stdio::piped(),
377            Self::Null => Stdio::null(),
378        }
379    }
380}
381
382/// Typed spawn description accepted by the blessed process boundary.
383#[cfg(feature = "async-process")]
384#[derive(Debug, Clone)]
385pub struct SpawnSpec {
386    program: OsString,
387    args: Vec<OsString>,
388    current_dir: Option<PathBuf>,
389    env: Vec<(OsString, OsString)>,
390    clear_env: bool,
391    stdin: StreamMode,
392    stdout: StreamMode,
393    stderr: StreamMode,
394    create_process_group: bool,
395    kill_when_owner_dies: bool,
396    nice: Option<i32>,
397    admission: Option<SpawnAdmission>,
398}
399
400#[cfg(feature = "async-process")]
401impl SpawnSpec {
402    /// Create a direct (non-shell) command description.
403    pub fn new(program: impl Into<OsString>) -> Self {
404        Self {
405            program: program.into(),
406            args: Vec::new(),
407            current_dir: None,
408            env: Vec::new(),
409            clear_env: false,
410            stdin: StreamMode::Inherit,
411            stdout: StreamMode::Inherit,
412            stderr: StreamMode::Inherit,
413            create_process_group: false,
414            kill_when_owner_dies: false,
415            nice: None,
416            admission: None,
417        }
418    }
419
420    /// Append one argument without requiring UTF-8.
421    pub fn arg(mut self, arg: impl Into<OsString>) -> Self {
422        self.args.push(arg.into());
423        self
424    }
425
426    /// Set the child working directory.
427    pub fn current_dir(mut self, path: impl Into<PathBuf>) -> Self {
428        self.current_dir = Some(path.into());
429        self
430    }
431
432    /// Add an environment override.
433    pub fn env(mut self, key: impl Into<OsString>, value: impl Into<OsString>) -> Self {
434        self.env.push((key.into(), value.into()));
435        self
436    }
437
438    /// Start with an empty inherited environment before applying overrides.
439    pub fn clear_env(mut self, clear: bool) -> Self {
440        self.clear_env = clear;
441        self
442    }
443
444    /// Configure child stdin.
445    pub fn stdin(mut self, mode: StreamMode) -> Self {
446        self.stdin = mode;
447        self
448    }
449
450    /// Configure child stdout.
451    pub fn stdout(mut self, mode: StreamMode) -> Self {
452        self.stdout = mode;
453        self
454    }
455
456    /// Configure child stderr.
457    pub fn stderr(mut self, mode: StreamMode) -> Self {
458        self.stderr = mode;
459        self
460    }
461
462    /// Put the child in its own process group.
463    ///
464    /// This is what makes a group-wide soft signal addressable at all:
465    /// [`PlatformEmergencySignal::terminate_group_soft`] is a no-op without
466    /// it, because on POSIX the negative-PID signal would otherwise reach the
467    /// caller's own group, and on Windows `GenerateConsoleCtrlEvent` only
468    /// routes to children spawned with `CREATE_NEW_PROCESS_GROUP`. It also
469    /// detaches the child from the parent's console Ctrl+C, so it is opt-in.
470    pub fn create_process_group(mut self, create: bool) -> Self {
471        self.create_process_group = create;
472        self
473    }
474
475    /// Kill this child when the spawning process exits unexpectedly.
476    ///
477    /// Linux uses `PR_SET_PDEATHSIG(SIGTERM)` plus a pre-exec hard-exit race
478    /// guard when the parent changed before that signal could be armed.
479    /// Windows assigns the child to a
480    /// process-wide kill-on-close Job Object. macOS forks a kqueue supervisor
481    /// before exec and reports spawn success only after its owner and child
482    /// watches are registered.
483    pub fn kill_when_owner_dies(mut self, kill: bool) -> Self {
484        self.kill_when_owner_dies = kill;
485        self
486    }
487
488    /// Apply the host's existing niceness policy at child creation.
489    ///
490    /// On Unix this is the requested `setpriority(PRIO_PROCESS)` niceness.
491    /// Windows maps the established niceness bands to process creation
492    /// priority classes; it is deliberately a coarse host mapping rather
493    /// than a claim that numeric nice values are portable.
494    pub fn nice(mut self, nice: Option<i32>) -> Self {
495        self.nice = nice;
496        self
497    }
498
499    /// Select portable scheduling intent at native process creation.
500    pub fn priority(self, priority: ProcessPriority) -> Self {
501        self.nice(priority.nice_value())
502    }
503
504    /// Request portable scheduling intent where host policy permits it.
505    ///
506    /// The current host launch boundary applies this at creation; platforms
507    /// that reject the requested class retain their native error behavior.
508    pub fn priority_best_effort(self, priority: ProcessPriority) -> Self {
509        self.priority(priority)
510    }
511
512    /// Acquire a caller-owned permit around the native spawn attempt.
513    pub fn spawn_admission(mut self, admission: SpawnAdmission) -> Self {
514        self.admission = Some(admission);
515        self
516    }
517
518    /// Spawn using the canonical asynchronous platform operation.
519    pub async fn spawn(self) -> io::Result<PlatformChild> {
520        let mut command = Command::new(&self.program);
521        command.args(&self.args);
522        if let Some(current_dir) = self.current_dir.as_deref() {
523            command.current_dir(current_dir);
524        }
525        if self.clear_env {
526            command.env_clear();
527        }
528        for (key, value) in &self.env {
529            command.env(key, value);
530        }
531        command
532            .stdin(self.stdin.apply())
533            .stdout(self.stdout.apply())
534            .stderr(self.stderr.apply());
535        platform_imp::configure_command(
536            &mut command,
537            self.create_process_group,
538            self.kill_when_owner_dies,
539            self.nice,
540        )?;
541
542        let mut spawn = || command.spawn();
543        let child = match self.admission.as_ref() {
544            Some(admission) => admission.run(spawn)?,
545            None => spawn()?,
546        };
547        platform_imp::after_spawn(&child, self.kill_when_owner_dies)?;
548        Ok(PlatformChild::new(child, self.create_process_group))
549    }
550}
551
552/// Owned child handle returned by [`SpawnSpec::spawn`].
553#[cfg(feature = "async-process")]
554pub struct PlatformChild {
555    child: Child,
556    stdin: Option<ChildStdin>,
557    stdout: Option<ChildStdout>,
558    stderr: Option<ChildStderr>,
559    signal: PlatformEmergencySignal,
560}
561
562#[cfg(feature = "async-process")]
563impl PlatformChild {
564    fn new(mut child: Child, own_process_group: bool) -> Self {
565        let signal = PlatformEmergencySignal {
566            identity: async_child_identity(&child),
567            own_process_group,
568            // The legacy AsyncProcess actor historically retained only this
569            // numeric child-group leader on macOS. Keep it launch-bound for
570            // that API's compatibility path; sessions deliberately never use
571            // it because their control capability promises identity safety.
572            legacy_group_pid: child.id(),
573        };
574        Self {
575            stdin: child.stdin.take(),
576            stdout: child.stdout.take(),
577            stderr: child.stderr.take(),
578            child,
579            signal,
580        }
581    }
582
583    /// Return the operating-system process identifier, if available.
584    pub fn id(&self) -> Option<u32> {
585        self.child.id()
586    }
587
588    /// Wait for completion without capturing output.
589    pub async fn wait(&mut self) -> io::Result<ExitStatus> {
590        self.child.wait().await
591    }
592
593    /// Terminate the child and wait for its exit.
594    pub async fn kill(&mut self) -> io::Result<()> {
595        self.child.kill().await
596    }
597
598    /// Capture piped stdout and stderr while waiting for the child.
599    pub async fn wait_with_output(self) -> io::Result<Output> {
600        let Self {
601            mut child,
602            stdin,
603            stdout,
604            stderr,
605            ..
606        } = self;
607        // Match Tokio's `Child::wait_with_output` contract: one-shot output
608        // closes an owned stdin pipe so a child waiting for EOF can finish.
609        drop(stdin);
610        let (status, stdout, stderr) = tokio::try_join!(
611            child.wait(),
612            read_owned_to_end(stdout),
613            read_owned_to_end(stderr),
614        )?;
615        Ok(Output {
616            status,
617            stdout,
618            stderr,
619        })
620    }
621
622    /// Write bytes to piped stdin and flush them.
623    pub async fn write_stdin(&mut self, bytes: &[u8]) -> io::Result<()> {
624        let stdin = self.stdin.as_mut().ok_or_else(stdin_not_piped)?;
625        stdin.write_all(bytes).await?;
626        stdin.flush().await
627    }
628
629    /// Close the piped stdin handle, delivering EOF to the child.
630    ///
631    /// This operation is idempotent. Closing an inherited or null stdin is
632    /// also a no-op because there is no owned pipe to close.
633    pub fn close_stdin(&mut self) {
634        drop(self.stdin.take());
635    }
636
637    /// Read all bytes from piped stdout without waiting for process exit.
638    pub async fn read_stdout_to_end(&mut self) -> io::Result<Vec<u8>> {
639        let stdout = self.stdout.as_mut().ok_or_else(stdout_not_piped)?;
640        let mut bytes = Vec::new();
641        stdout.read_to_end(&mut bytes).await?;
642        Ok(bytes)
643    }
644
645    /// Read all bytes from piped stderr without waiting for process exit.
646    pub async fn read_stderr_to_end(&mut self) -> io::Result<Vec<u8>> {
647        let stderr = self.stderr.as_mut().ok_or_else(stderr_not_piped)?;
648        let mut bytes = Vec::new();
649        stderr.read_to_end(&mut bytes).await?;
650        Ok(bytes)
651    }
652
653    /// Split this child into sealed actor capabilities.
654    ///
655    /// The lifecycle wait handle, emergency termination handle, input pipe,
656    /// and output readers are deliberately separate so the actor can keep
657    /// accepting control commands while an asynchronous exit wait is pending.
658    pub fn into_actor_parts(
659        self,
660    ) -> (
661        PlatformLifecycle,
662        PlatformEmergencySignal,
663        Option<PlatformStdin>,
664        Option<PlatformOutput>,
665        Option<PlatformOutput>,
666    ) {
667        (
668            PlatformLifecycle { child: self.child },
669            self.signal,
670            self.stdin.map(|stdin| PlatformStdin { stdin }),
671            self.stdout.map(PlatformOutput::stdout),
672            self.stderr.map(PlatformOutput::stderr),
673        )
674    }
675}
676
677/// Opaque exit-wait capability owned by a process actor.
678#[cfg(feature = "async-process")]
679pub struct PlatformLifecycle {
680    child: Child,
681}
682
683#[cfg(feature = "async-process")]
684impl PlatformLifecycle {
685    /// Wait asynchronously for the child to exit.
686    pub async fn wait(&mut self) -> io::Result<ExitStatus> {
687        self.child.wait().await
688    }
689
690    /// Request direct-child termination through the still-owned child handle
691    /// without waiting for its reaping result.
692    ///
693    /// This is the identity-safe fallback when a host cannot provide a
694    /// separately usable launch-bound signal capability (for example a Linux
695    /// kernel without pidfds). The actor continues to own this lifecycle
696    /// handle and performs the eventual reap itself.
697    pub fn start_kill(&mut self) -> io::Result<()> {
698        self.child.start_kill()
699    }
700}
701
702/// Opaque, non-reap-capable emergency termination capability.
703///
704/// It can be used while the actor has a pending wait on
705/// [`PlatformLifecycle`], but it cannot observe or consume the exit result.
706#[cfg(feature = "async-process")]
707pub struct PlatformEmergencySignal {
708    identity: Option<AsyncChildIdentity>,
709    own_process_group: bool,
710    legacy_group_pid: Option<u32>,
711}
712
713#[cfg(feature = "async-process")]
714impl PlatformEmergencySignal {
715    /// Request immediate termination without waiting for process reaping.
716    pub fn kill(&self) -> io::Result<()> {
717        let Some(identity) = self.identity.as_ref() else {
718            return Err(signal_target_unavailable());
719        };
720        signal_async_child(identity)
721    }
722
723    /// Ask the child's whole process group to shut down gracefully.
724    ///
725    /// Returns `Ok(false)` when the child was not spawned with
726    /// [`SpawnSpec::create_process_group`]: there is no group to address, and
727    /// signalling anyway would hit the caller's own group on POSIX or the
728    /// caller's console on Windows. A missing or mismatched launch identity
729    /// instead reports an unavailable target; it never falls back to a
730    /// numeric group identifier that might have been reused.
731    pub fn terminate_group_soft(&self) -> io::Result<bool> {
732        if !self.own_process_group {
733            return Ok(false);
734        }
735        let Some(identity) = self.identity.as_ref() else {
736            return Err(signal_target_unavailable());
737        };
738        signal_async_child_group(identity).map(|()| true)
739    }
740
741    /// Legacy AsyncProcess-only graceful group termination.
742    ///
743    /// Most hosts retain an identity-safe asynchronous signal capability. On
744    /// macOS there is no pidfd-equivalent for the Tokio child path, while the
745    /// pre-session AsyncProcess contract historically sent SIGTERM to the
746    /// launch child's numeric process group. Preserve that established
747    /// best-effort behavior only for the legacy actor; sessions keep using
748    /// [`Self::terminate_group_soft`] and therefore remain identity-safe.
749    pub fn terminate_group_soft_legacy(&self) -> io::Result<bool> {
750        if !self.own_process_group {
751            return Ok(false);
752        }
753        if let Some(identity) = self.identity.as_ref() {
754            return signal_async_child_group(identity).map(|()| true);
755        }
756        let pid = self
757            .legacy_group_pid
758            .ok_or_else(signal_target_unavailable)?;
759        crate::platform::process::soft_terminate_process_group(pid).map(|()| true)
760    }
761
762    /// Return direct-child CPU time when this host can still verify the
763    /// launch identity. Unsupported hosts and already-reused identities are
764    /// reported as `None`, never as a PID-only best effort.
765    pub fn cpu_time(&self) -> io::Result<Option<std::time::Duration>> {
766        self.identity
767            .as_ref()
768            .map_or(Ok(None), async_child_cpu_time)
769    }
770}
771
772#[cfg(feature = "async-process")]
773fn signal_target_unavailable() -> io::Error {
774    io::Error::new(
775        io::ErrorKind::BrokenPipe,
776        "child process launch identity is no longer available",
777    )
778}
779
780/// Opaque piped stdin capability owned by a process actor.
781#[cfg(feature = "async-process")]
782pub struct PlatformStdin {
783    stdin: ChildStdin,
784}
785
786#[cfg(feature = "async-process")]
787impl PlatformStdin {
788    /// Write and flush bytes to the child stdin pipe.
789    pub async fn write(&mut self, bytes: &[u8]) -> io::Result<()> {
790        self.stdin.write_all(bytes).await?;
791        self.stdin.flush().await
792    }
793}
794
795/// Opaque stdout or stderr reader owned by a process actor.
796#[cfg(feature = "async-process")]
797pub struct PlatformOutput {
798    reader: OutputReader,
799    read_pending: bool,
800}
801
802#[cfg(feature = "async-process")]
803enum OutputReader {
804    Stdout(ChildStdout),
805    Stderr(ChildStderr),
806}
807
808#[cfg(feature = "async-process")]
809impl PlatformOutput {
810    fn stdout(stdout: ChildStdout) -> Self {
811        Self {
812            reader: OutputReader::Stdout(stdout),
813            read_pending: false,
814        }
815    }
816
817    fn stderr(stderr: ChildStderr) -> Self {
818        Self {
819            reader: OutputReader::Stderr(stderr),
820            read_pending: false,
821        }
822    }
823
824    /// Drain this output endpoint to EOF without blocking a runtime worker.
825    pub async fn read_to_end(self) -> io::Result<Vec<u8>> {
826        match self.reader {
827            OutputReader::Stdout(stdout) => read_owned_to_end(Some(stdout)).await,
828            OutputReader::Stderr(stderr) => read_owned_to_end(Some(stderr)).await,
829        }
830    }
831
832    /// Read the next asynchronous chunk from this output endpoint.
833    ///
834    /// The caller owns the buffer and therefore controls the amount of data
835    /// retained at each read. EOF is reported as `Ok(0)`.
836    pub async fn read_chunk(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
837        self.read_pending = true;
838        let result = match &mut self.reader {
839            OutputReader::Stdout(stdout) => stdout.read(buffer).await,
840            OutputReader::Stderr(stderr) => stderr.read(buffer).await,
841        };
842        self.read_pending = false;
843        result
844    }
845
846    /// Abandon output and await destruction of this reader's pending I/O.
847    ///
848    /// This is not direct-child reaping or lossless EOF. The owning actor must
849    /// retain this future until completion: dropping it is not an acknowledgement
850    /// that platform I/O storage has been released. Cancellation requests alone
851    /// do not count as completion, and an unresponsive driver may keep it pending.
852    pub async fn shutdown(self) -> io::Result<()> {
853        match self.reader {
854            OutputReader::Stdout(stdout) => {
855                platform_imp::shutdown_output_reader(stdout, self.read_pending).await
856            }
857            OutputReader::Stderr(stderr) => {
858                platform_imp::shutdown_output_reader(stderr, self.read_pending).await
859            }
860        }
861    }
862}
863
864#[cfg(all(test, feature = "async-process"))]
865mod output_shutdown_tests {
866    use super::*;
867    use std::time::Duration;
868
869    #[test]
870    fn silent_output_fixture() {
871        if std::env::var_os("RUNNING_PROCESS_OUTPUT_SHUTDOWN_FIXTURE").is_some() {
872            std::thread::sleep(Duration::from_secs(30));
873        }
874    }
875
876    #[tokio::test]
877    async fn shutdown_finishes_a_cancelled_pending_read_before_child_exit() {
878        let child = SpawnSpec::new(std::env::current_exe().expect("test executable"))
879            .arg("--exact")
880            .arg("output_shutdown_tests::silent_output_fixture")
881            .env("RUNNING_PROCESS_OUTPUT_SHUTDOWN_FIXTURE", "1")
882            .stdin(StreamMode::Null)
883            .stdout(StreamMode::Piped)
884            .stderr(StreamMode::Null)
885            .spawn()
886            .await
887            .expect("spawn silent fixture");
888        let (mut lifecycle, _, _, stdout, _) = child.into_actor_parts();
889        let mut stdout = stdout.expect("stdout pipe");
890        // Drain the libtest header, then establish a genuinely pending read.
891        let mut bytes = [0; 1024];
892        loop {
893            match tokio::time::timeout(Duration::from_millis(20), stdout.read_chunk(&mut bytes))
894                .await
895            {
896                Ok(Ok(0)) => panic!("fixture exited before pending read"),
897                Ok(Ok(_)) => {}
898                Ok(Err(error)) => panic!("fixture read failed: {error}"),
899                Err(_) => break,
900            }
901        }
902        let shutdown = tokio::time::timeout(Duration::from_secs(2), stdout.shutdown()).await;
903        // Reap even if the assertion fails; a passing shutdown cannot depend
904        // on this kill, which occurs only after the acknowledgement deadline.
905        lifecycle.start_kill().expect("kill fixture");
906        lifecycle.wait().await.expect("reap fixture");
907        shutdown
908            .expect("pending read shutdown must finish")
909            .expect("output shutdown");
910    }
911}
912
913#[cfg(feature = "async-process")]
914fn stdin_not_piped() -> io::Error {
915    io::Error::new(io::ErrorKind::BrokenPipe, "child stdin is not piped")
916}
917
918#[cfg(feature = "async-process")]
919fn stdout_not_piped() -> io::Error {
920    io::Error::new(io::ErrorKind::BrokenPipe, "child stdout is not piped")
921}
922
923#[cfg(feature = "async-process")]
924fn stderr_not_piped() -> io::Error {
925    io::Error::new(io::ErrorKind::BrokenPipe, "child stderr is not piped")
926}
927
928#[cfg(feature = "async-process")]
929async fn read_owned_to_end<R>(reader: Option<R>) -> io::Result<Vec<u8>>
930where
931    R: AsyncRead + Unpin,
932{
933    let Some(mut reader) = reader else {
934        return Ok(Vec::new());
935    };
936    let mut bytes = Vec::new();
937    reader.read_to_end(&mut bytes).await?;
938    Ok(bytes)
939}
940
941/// Build a shell command using the host platform's supported shell.
942#[cfg(feature = "async-process")]
943pub fn shell_spec(command: impl AsRef<OsStr>) -> SpawnSpec {
944    platform_imp::shell_spec(command.as_ref())
945}
946
947#[cfg(all(test, feature = "async-process"))]
948mod tests {
949    use super::{shell_spec, SpawnSpec, StreamMode};
950
951    fn fixture_command() -> SpawnSpec {
952        #[cfg(windows)]
953        {
954            shell_spec("echo async-platform-internal")
955        }
956        #[cfg(not(windows))]
957        {
958            shell_spec("printf async-platform-internal")
959        }
960    }
961
962    #[tokio::test]
963    async fn blessed_spawn_captures_output_without_sync_wait() {
964        let output = fixture_command()
965            .stdout(StreamMode::Piped)
966            .stderr(StreamMode::Piped)
967            .spawn()
968            .await
969            .expect("spawn")
970            .wait_with_output()
971            .await
972            .expect("wait with output");
973
974        assert!(output.status.success());
975        let expected = if cfg!(windows) {
976            b"async-platform-internal\r\n".as_slice()
977        } else {
978            b"async-platform-internal".as_slice()
979        };
980        assert_eq!(output.stdout, expected);
981        assert!(output.stderr.is_empty());
982    }
983
984    #[tokio::test]
985    async fn blessed_spawn_reports_missing_program() {
986        let result = SpawnSpec::new("running-process-program-that-does-not-exist")
987            .spawn()
988            .await;
989        assert!(result.is_err());
990    }
991
992    #[tokio::test]
993    async fn one_shot_output_closes_owned_stdin() {
994        #[cfg(windows)]
995        let spec = shell_spec("more > nul & echo done");
996        #[cfg(not(windows))]
997        let spec = shell_spec("cat > /dev/null; printf done");
998
999        let output = tokio::time::timeout(
1000            std::time::Duration::from_secs(2),
1001            spec.stdin(StreamMode::Piped)
1002                .stdout(StreamMode::Piped)
1003                .stderr(StreamMode::Piped)
1004                .spawn()
1005                .await
1006                .expect("spawn")
1007                .wait_with_output(),
1008        )
1009        .await
1010        .expect("stdin is closed for one-shot output")
1011        .expect("output succeeds");
1012
1013        let expected = if cfg!(windows) {
1014            b"done\r\n".as_slice()
1015        } else {
1016            b"done".as_slice()
1017        };
1018        assert_eq!(output.stdout, expected);
1019    }
1020}