Skip to main content

running_process/
lib.rs

1//! Cross-platform process execution, process-tree control, PTY handling, and
2//! broker integration primitives.
3//!
4//! The crate exposes a synchronous process API through [`NativeProcess`], a
5//! contained process-group helper through [`ContainedProcessGroup`], low-level
6//! spawn helpers through [`spawn()`] and [`spawn_daemon`], and optional
7//! daemon/broker modules behind feature flags.
8// #1101: environment reads go through declared variables; see the
9// `running_process_env_direct` Dylint lint.
10#![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
25/// Explicit foreground commands preserving caller-controlled native launch state.
26pub 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// Phase 1 of #221: process-observation capability model + portable
49// lifecycle baseline. Core-feature-clean (std-only: mpsc + SystemTime),
50// so the started/exited baseline is available to the base library
51// without pulling in the daemon runtime.
52/// Frozen v1 daemon manifest and service-definition registration substrate.
53///
54/// This direct persistence surface owns only the v1 registration records,
55/// SHA-256 seal/verify rules, host stamp, validated paths, and private-file
56/// behavior. It deliberately does not select an IPC endpoint, broker client,
57/// daemon runtime, identity probe, or async runtime.
58#[cfg(feature = "daemon-registration")]
59pub mod daemon_registration;
60/// Frozen v1 semantic registration compatibility contract.
61#[cfg(feature = "daemon-registration")]
62pub mod daemon_registration_compat;
63/// Frozen v2 service-definition registration writer substrate.
64///
65/// This direct persistence surface owns the established `.servicedef.v2`
66/// layout, generated service definition, validation, and owner-private file
67/// behavior. It deliberately does not select v2 manifests, a loader, broker
68/// negotiation, endpoint transport, identity, or an async runtime.
69#[cfg(feature = "daemon-registration-v2")]
70pub mod daemon_registration_v2;
71/// Limited shared-broker v2 registration compatibility contract.
72#[cfg(feature = "daemon-registration-v2")]
73pub mod daemon_registration_v2_compat;
74// The two registration writer features share only the small path, name, error,
75// and owner-private-directory substrate. Keeping it separate from either
76// public module prevents v2 persistence from selecting v1's SHA-256 manifest
77// support, while retaining exact v1 type identity through re-exports.
78/// Canonical semantic v1 frame compatibility contract, retaining raw values.
79#[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/// Frozen v1 `Frame` envelope codec and consumer-protocol registry.
84///
85/// This direct, transport-free surface is available without broker IPC,
86/// daemon identity, hashing, or an async runtime. Broad broker paths
87/// re-export these exact items for compatibility.
88#[cfg(feature = "frame-v1-codec")]
89pub mod frame_v1;
90// Host facts are shared by the direct identity probe and persisted v1
91// registration. The implementation is deliberately private; registration
92// exposes its stable public host-identity path from `daemon_registration`.
93#[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// The IPC client owns the generated protocol dependency.  Keeping code
101// generation in that optional package means process-only consumers do not
102// compile broker schemas or their build dependencies (#1144).
103#[cfg(feature = "client")]
104/// Prost-generated daemon protocol types used by the client transport.
105pub mod proto {
106    /// Generated Rust bindings for the `running_process.daemon.v1` protobuf package.
107    pub use running_process_protocol::daemon;
108}
109
110#[cfg(feature = "client")]
111pub mod client;
112
113// Phase 0 of #228: v1 broker module — prost-generated wire types from
114// `proto/broker_v1_*.proto`. The broad broker remains a `client` API. The
115// narrow identity substrate compiles this module privately so its direct
116// facade can preserve the frozen v1 probe/frame bytes without exposing broker
117// ownership, configuration, or client APIs.
118#[cfg(feature = "client")]
119pub mod broker;
120// The direct facade imports a deliberately small subset of the legacy
121// namespace while its compatibility re-exports remain available for type
122// identity.  The remaining client-only paths are intentionally dormant here.
123#[cfg(all(feature = "backend-identity", not(feature = "client")))]
124#[allow(dead_code, unused_imports)]
125mod broker;
126
127/// Direct daemon-identity substrate for an existing endpoint.
128///
129/// This is intentionally a small facade over the frozen v1 identity probe,
130/// sidecar, and sans-I/O endpoint mux. It does not adopt the broker client or
131/// daemon runtime, and it leaves endpoint naming and application payloads to
132/// the caller.
133#[cfg(feature = "backend-identity")]
134pub mod backend_identity;
135
136// #891: content-hash primitive (`blake3_file`) for dev daemon-identity
137// isolation. The direct identity facade needs it internally, but it remains a
138// public client-only utility rather than widening the direct facade.
139#[cfg(feature = "client")]
140pub mod content_hash;
141#[cfg(all(feature = "backend-identity", not(feature = "client")))]
142mod content_hash;
143
144/// Probe client facade (#633). Gated on the `probe` feature so a build
145/// without it contains none of this code.
146#[cfg(feature = "probe")]
147pub mod probe;
148
149// Phase 1 of #228 (issue #230): maintenance subcommands exposed via
150// the `runpm` CLI. Currently just `release-handles` — a cross-platform
151// scaffold for the Windows worktree-teardown handle-race fix
152// (soldr#710). Gated behind `feature = "client"` because the CLI that
153// drives it is.
154#[cfg(feature = "client")]
155pub mod maintenance;
156
157#[cfg(feature = "client")]
158pub mod cleanup;
159
160// Phase 4 of #222 (#427): per-OS boot autostart for the runpm daemon.
161// Gated behind `feature = "client"` because the only consumer is the
162// `runpm` CLI binary, which is itself client-gated.
163#[cfg(feature = "client")]
164pub mod boot_autostart;
165
166// Phase 5 of #222 (#428): `runpm.toml` parser used by the `runpm` CLI
167// to batch-start `[[app]]` entries. Lives in the library (not under
168// `src/bin/`) so the integration test in `tests/runpm/runpm_toml_config.rs`
169// can drive the same code path the binary uses.
170#[cfg(feature = "client")]
171pub mod runpm_config;
172
173// #415: consumer-consumable conformance test kit. Gated behind the
174// off-by-default `test-support` cargo feature (which implies `client`)
175// so the helpers ship in the published crate but only compile when a
176// consumer opts in as a dev-dependency.
177#[cfg(feature = "test-support")]
178pub mod test_support;
179
180// Lightweight tee sink primitives for callers that want transcript/log
181// fan-out without pulling in the full daemon runtime.
182//
183// The file lives under `daemon/` because that is who else uses it, and the
184// `daemon` feature loads it there as `daemon::telemetry`. Declaring it as a
185// module here too would load one file as two modules -- two copies of every
186// type, which are then not the same type -- so when both features are on this
187// re-exports the daemon's module instead of declaring a second one.
188#[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/// `telemetry` and `daemon::telemetry` must name one module, not two copies.
196///
197/// A `#[path]` module declaration alongside the daemon's own would compile --
198/// that was the bug -- but it would mint a second, incompatible set of types
199/// from the same source file, so a `TeeHandle` obtained through one path
200/// could not be passed to a function expecting the other. This conversion is
201/// the identity only while both paths resolve to the same item; if the
202/// duplicate declaration ever comes back, it stops compiling here rather than
203/// at whichever caller first tried to mix the two.
204#[cfg(all(feature = "telemetry", feature = "daemon"))]
205const _: fn(crate::telemetry::TeeHandle) -> daemon::telemetry::TeeHandle = |handle| handle;
206
207// Wave 5 of #165: daemon runtime absorbed from `running-process-daemon`.
208// Heavy deps (tokio, sqlite, etc.) gated behind `feature = "daemon"`.
209#[cfg(feature = "daemon")]
210/// Daemon runtime APIs and helpers enabled by the `daemon` feature.
211pub mod daemon;
212// `kill_tree` is established 4.x containment surface and remains available to
213// `default-features = false` callers. Its sysinfo-backed platform primitive is
214// the explicit Phase 0.5 compatibility exception; public inspection APIs stay
215// behind `process-inspection`.
216#[cfg(feature = "independent-spawn")]
217pub mod independent_spawn;
218pub mod process_tree;
219#[cfg(feature = "pty")]
220/// PTY-backed process APIs.
221pub 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// #891: content-hash primitive for dev daemon-identity isolation.
248#[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/// Executable naming and image-relative discovery, for binaries in this
266/// workspace that must name a sibling program without spelling it per host.
267#[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;
271/// Canonical native process-inspection errors, preserving their host detail.
272pub use running_process_platform_internal::platform::process::{
273    ProcessInspectError, ProcessInspectErrorKind,
274};
275/// Resolve the current executable image for a live PID.
276pub use running_process_platform_internal::process_executable_path;
277/// Compare executable path spellings using the host-native policy.
278pub use running_process_platform_internal::process_same_executable_path;
279/// Retained native process-liveness observation.
280pub 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};
315/// Convert a native process exit status to the portable integer convention.
316pub 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]
325/// Create a scoped Rust debug trace label for the current function body.
326macro_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    /// Byte-exact stream chunks. Unlike the logical line queues these retain
342    /// delimiters, unterminated tails, and non-UTF-8 bytes; callers consume
343    /// them with `drain_stream_raw`.
344    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
355/// Sentinel value for returncode atomic: process has not exited yet.
356const 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    /// Atomic exit code. `RETURNCODE_NOT_SET` means "not exited yet".
365    /// Updated by the child actor — reading is lock-free.
366    returncode: AtomicI64,
367    /// Phase 1 of #221: optional lifecycle-event emitter. `None` means
368    /// observation is off (the off-by-default path), so the lifecycle
369    /// hooks are inert. When `Some`, `started` is emitted once at spawn
370    /// and `exited` exactly once on the first returncode transition.
371    observer: Option<ObserverEmitter>,
372    /// Guards against emitting more than one `exited` event when several
373    /// code paths (the child actor serving its tick, `poll`, `kill`) race to record the exit.
374    observer_exit_emitted: AtomicBool,
375    /// #850: exit publication for waiters that run on the actor runtime.
376    /// Mirrors `returncode`; every write goes through [`Self::record_exit`].
377    exit_code: tokio::sync::watch::Sender<Option<i32>>,
378}
379
380/// The child owned by a started [`NativeProcess`]: the platform's non-Tokio
381/// backend (a std or exact-trace child, observed by `try_wait` polling from
382/// the actor runtime), which also owns the Windows per-spawn Job Object and
383/// the capture cancellation for its pipes (#850).
384type 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    /// Publish the exit code to every observer: lock-free readers of
439    /// `returncode`, condvar waiters, and actor-runtime waiters.
440    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    /// Emit the lifecycle `exited` event exactly once, regardless of which
447    /// code path first observes the exit. No-op when observation is off.
448    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
462/// A cross-platform child process with optional output capture.
463///
464/// `NativeProcess` wraps [`std::process::Command`] with the crate's
465/// process-tree containment, capture draining, timeout, and terminal-control
466/// behavior. Methods are synchronous and are safe to call from ordinary
467/// blocking code.
468pub struct NativeProcess {
469    config: ProcessConfig,
470    command_override: Mutex<Option<Command>>,
471    /// Set once by `start`; the child itself lives in its actor.
472    child: std::sync::OnceLock<child_actor::ChildHandle>,
473    /// Serializes concurrent `start` calls; never held by any other path.
474    start_gate: Mutex<()>,
475    stdin: Mutex<Option<ChildStdin>>,
476    shared: Arc<SharedState>,
477    process_watch: Option<Arc<ProcessWatchEmitter>>,
478    // This remains a constructor-only policy for the bounded std::Command
479    // entrypoint. General NativeProcess callers keep their established
480    // platform policy surface.
481    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    /// Create a process wrapper from a [`ProcessConfig`].
490    ///
491    /// The child is not spawned until [`Self::start`] is called. Process
492    /// observation is **off by default**: no lifecycle events are emitted
493    /// unless [`Self::with_observer`] is used instead.
494    pub fn new(config: ProcessConfig) -> Self {
495        Self::new_with_options(config, None, None, None, None)
496    }
497
498    /// Create a process wrapper with process observation enabled (Phase 1
499    /// of #221).
500    ///
501    /// Returns the wrapper paired with an [`ObserverSubscriber`] that
502    /// receives a [`started`](crate::ObserverEventKind::Started) event when
503    /// [`Self::start`] spawns the child and exactly one
504    /// [`exited`](crate::ObserverEventKind::Exited) event when the child is
505    /// reaped — for the categories the `config` requests that are actually
506    /// `Supported` (only [`Lifecycle`](crate::EventCategory::Lifecycle) in
507    /// Phase 1; see [`ObserverCapabilities::negotiate`](crate::ObserverCapabilities::negotiate)).
508    ///
509    /// The emitter never blocks on a slow or dropped subscriber.
510    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    /// Create a process wrapper from a caller-configured
520    /// [`std::process::Command`] with process observation enabled.
521    ///
522    /// [`ProcessConfig`]'s declarative surface deliberately cannot represent
523    /// everything a `Command` can carry — `env_remove` scrubs of inherited
524    /// variables, non-Unicode (`OsString`) argv/env values, a pre-resolved
525    /// working directory — so a caller that already owns a fully configured
526    /// `Command` (a compiler front door wrapping cargo, zackees/soldr#2546)
527    /// would otherwise have to lossily re-encode it. This pairs the observer
528    /// machinery with the same command-override seam the capture-limit
529    /// constructors use: `command` is spawned verbatim, while `config` still
530    /// governs stdio routing, capture, containment, and limits (its
531    /// `command` / `cwd` / `env` fields are ignored in favor of the
532    /// override, matching `build_command`).
533    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    /// Create a process with launched-tree watch matching configured before
544    /// spawn. Exact tracing, when selected, owns the launch-time wait events.
545    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    /// Describe exact launched-tree observation support on this host.
556    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    // Preserve a stable Rust frame here in release user dumps.
600    #[inline(never)]
601    /// Spawn the configured child process.
602    ///
603    /// Returns [`ProcessError::AlreadyStarted`] if the same wrapper already
604    /// owns a running child.
605    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        // Phase 1 of #221: emit the lifecycle `started` event. No-op when
663        // observation is off (the common, off-by-default path).
664        if let Some(emitter) = self.shared.observer.as_ref() {
665            emitter.emit_started(child.id());
666        }
667        // #539 slice 2: when the observer requests EventCategory::Process,
668        // associate an IOCP with the per-spawn Job Object so a pump thread
669        // can forward descendant lifecycle events. The Lifecycle category
670        // is still served by emit_started / emit_exited above and below.
671        // Hosts without Job Objects answer `Unsupported`, which means there is
672        // nothing to contain; an exact-trace child has no standard handle to
673        // assign.
674        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    /// Write bytes to the child's stdin and then close stdin.
744    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    /// Write to the child's stdin without closing it afterwards, so the
763    /// caller can issue additional writes. Used by interactive
764    /// pipe-backed sessions (#130 milestone 3) where the daemon keeps
765    /// stdin open across multiple client input frames.
766    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    /// Explicitly close the child's stdin (signals EOF to the child).
784    /// Idempotent: returns Ok if stdin was already closed.
785    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    /// Check whether the child has exited without blocking.
794    ///
795    /// Returns `Ok(None)` while the process is still running.
796    pub fn poll(&self) -> Result<Option<i32>, ProcessError> {
797        // Fast path: check atomic set by the child actor.
798        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        // The actor publishes the exit (returncode, watch, lifecycle event)
805        // when this is the check that finds it.
806        child.try_wait().map_err(ProcessError::Io)
807    }
808
809    // Preserve a stable Rust frame here in release user dumps.
810    #[inline(never)]
811    /// Wait for the child to exit.
812    ///
813    /// When `timeout` is `Some`, returns [`ProcessError::Timeout`] if the
814    /// child does not exit before the duration elapses.
815    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        // Fast path: already exited.
825        if let Some(code) = self.returncode() {
826            self.finish_capture_drain();
827            return Ok(code);
828        }
829        // A short timed wait must not depend on the child actor's tick. That
830        // task runs on the runtime's timer, whose granularity is coarse on some
831        // hosts (about 15 ms on Windows), so a child that has already exited can
832        // stay unobserved for longer than a caller's grace period -- which made
833        // `stream_iter` emit a second terminal event when its 10 ms grace lapsed
834        // right after EOF. For the first stretch of a timed wait, check the
835        // child directly on the calling thread, where a sleep is precise.
836        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                // A failed check is treated as "still running": the lifecycle
842                // tick, or the timeout below, decides what happens next.
843                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        // #850: the exit is published by the child actor on the actor
867        // runtime. `block_on_anywhere` is safe from a Tokio worker too, so a
868        // sync caller inside async code keeps working rather than erroring.
869        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            // The sender lives in `self.shared`, so it cannot close while
885            // `self` is borrowed; treat that impossibility as not running.
886            Some(None) => Err(ProcessError::NotRunning),
887            Some(Some(code)) => {
888                self.finish_capture_drain();
889                Ok(code)
890            }
891        }
892    }
893
894    // Preserve a stable Rust frame here in release user dumps.
895    #[inline(never)]
896    /// Forcefully terminate the child process.
897    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        // The actor checks for an exit first and otherwise signals the child
906        // (group-wide when it leads its own group and the host can).
907        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        // Wake capture readers immediately after the kill. In particular, this
913        // prevents a surviving pipe-owning descendant (FastLED Bug B: `uv`
914        // spawns a `python` grandchild that inherits the pipe and outlives it)
915        // from extending the bounded reap window, and wakes a blocked
916        // `read()` in microseconds rather than at the drain deadline.
917        self.cancel_capture_io();
918        // The actor publishes the exit and the lifecycle `exited` event
919        // (Phase 1 of #221: a killed child still produces one); all this has
920        // to do is give it the bounded window to observe the reap.
921        if already_reaped.is_none() {
922            self.await_exit_until(deadline);
923        }
924        // Synchronize with the per-stream reader threads so that by the time
925        // kill() returns, the capture queues have flipped from "blocked on
926        // read" to "closed" and downstream pollers (e.g. take_combined_line)
927        // observe EOS instead of timeout. The deadline remains a safety net
928        // if the platform wake mechanism does not fire.
929        public_symbols::rp_native_process_wait_for_capture_completion_with_deadline_public(
930            self, deadline,
931        );
932        Ok(())
933    }
934
935    /// Wait for the child actor to publish the exit, up to `deadline`. A last
936    /// `try_wait` through the actor covers a lifecycle tick that has already
937    /// stopped (for instance after a terminal `try_wait` error).
938    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    /// Terminate the child process.
951    ///
952    /// This currently uses the same hard-kill path as [`Self::kill`].
953    pub fn terminate(&self) -> Result<(), ProcessError> {
954        self.kill()
955    }
956
957    /// Send the OS-appropriate soft termination signal to the child's
958    /// process group (POSIX: SIGTERM to `-pid`; Windows: Ctrl+Break).
959    ///
960    /// Requires `ProcessConfig.create_process_group=true` on POSIX so
961    /// that `-pid` resolves to the child's own group. With the default
962    /// `create_process_group=false`, the kill would walk back to the
963    /// caller's group; the method silently no-ops in that case to avoid
964    /// signaling the wrong tree.
965    ///
966    /// Used by the daemon-side pipe sessions (#130 M4 follow-up) so
967    /// that `TerminationOutcome::SoftExit` becomes meaningful on POSIX.
968    pub fn terminate_group_soft(&self) -> Result<(), ProcessError> {
969        if !self.config.create_process_group {
970            // A group signal would otherwise reach the caller's own group.
971            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    // Preserve a stable Rust frame here in release user dumps.
979    #[inline(never)]
980    /// Close the process wrapper by terminating the child when it is running.
981    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    /// Return the child process id when the wrapper currently owns a child.
1002    pub fn pid(&self) -> Option<u32> {
1003        self.child.get().map(child_actor::ChildHandle::pid)
1004    }
1005
1006    /// Return the cached exit code when the child has exited.
1007    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    /// Return whether captured output is queued for one stream.
1017    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    /// Return whether captured combined output is queued.
1029    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    /// Drain and return all queued output for one stream.
1035    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    /// Consume and return all byte-exact output currently captured for one
1048    /// stream.
1049    ///
1050    /// This is independent of the logical-line queues used by
1051    /// [`Self::read_stream`] and [`Self::drain_stream`]. It preserves CRLF/LF
1052    /// delimiters, unterminated tails, and non-UTF-8 bytes exactly as accepted
1053    /// from the pipe reader. Calling it after EOF returns an empty vector once
1054    /// the stream has been drained.
1055    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    /// Drain and return all queued combined output events.
1081    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    /// Read the next captured chunk from one stream.
1087    ///
1088    /// Returns [`ReadStatus::Timeout`] when `timeout` elapses before output or
1089    /// EOF is observed.
1090    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    // Preserve a stable Rust frame here in release user dumps.
1154    #[inline(never)]
1155    /// Read the next captured combined stream event.
1156    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    /// Return the retained stdout history.
1202    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    /// Return the retained stderr history.
1219    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    /// Return the retained combined stdout/stderr event history.
1242    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    /// Return the retained byte count for one captured stream.
1254    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    /// Return the retained byte count for combined captured output.
1266    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    /// Clear retained output history for one stream and return freed bytes.
1275    ///
1276    /// This releases both the logical-line history and the byte-exact queue
1277    /// behind [`Self::drain_stream_raw`], so it stays the single memory-release
1278    /// valve for a captured stream. The returned count is the logical-line
1279    /// history only, unchanged. A caller consuming byte-exact output should
1280    /// drain it with [`Self::drain_stream_raw`] — which frees it too — rather
1281    /// than interleaving this call.
1282    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    /// Clear retained combined output history and return freed bytes.
1308    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                // The child's PATH to set last, when an APE launch puts its
1326                // loader's directory first on it.
1327                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                        // An APE image runs through its planned loader on a
1332                        // host that cannot exec it (see `crate::ape`).
1333                        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            // Clear the parent-side pipe-handle slot under its mutex
1427            // before dropping the reader. After this returns,
1428            // `kill_impl` can no longer try to `CancelIoEx` on us, so
1429            // it's safe for `reader`'s drop to close the HANDLE.
1430            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    /// Cancel pending capture reads so reader threads return immediately.
1462    /// Used by `kill_impl` to break the grandchild-orphan deadlock without
1463    /// waiting on `wait_for_capture_completion_with_deadline`'s safety-net.
1464    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    /// Bounded capture drain for the natural-exit and `close` paths
1477    /// (issue #590, cluster A). Waits at most `kill_drain_deadline` for the
1478    /// reader threads to flip the closed flags, force-setting them on
1479    /// timeout so `wait()`/`close()` return in bounded time instead of
1480    /// wedging in the previously-unbounded `wait_for_capture_completion`.
1481    /// Unlike `kill_impl` the reader is not cancelled up front — a
1482    /// short-lived grandchild's output is allowed to drain within the
1483    /// grace window — but if the window elapses with the pipe still held
1484    /// open the reader is cancelled to release the leaked thread.
1485    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    /// Returns `true` if the reader threads flipped both closed flags on their
1497    /// own before `deadline`, `false` if the deadline forced completion.
1498    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
1530/// How long a timed `wait` checks the child directly before falling back to
1531/// the child actor's lifecycle tick. Long enough to cover the coarsest timer tick a
1532/// host has, short enough that the polling never becomes the steady-state cost.
1533const SHORT_WAIT_DIRECT_POLL: Duration = Duration::from_millis(50);
1534
1535/// Cancel any pending blocking `read()` on the parent-side capture pipes
1536/// so the reader threads' `read()` calls return `ERROR_OPERATION_ABORTED`
1537/// immediately. Shared by `kill_impl`, `poll`, and the natural-exit
1538/// child actor (issue #590) — anywhere the child is observed to exit
1539/// while a grandchild may still hold the pipe open.
1540/// Wait until both capture streams report closed or `deadline` elapses.
1541/// On deadline, force-set the closed flags (and notify all waiters) so
1542/// downstream pollers observe EOF instead of blocking forever. Returns
1543/// `true` if the reader threads flipped the flags on their own, `false`
1544/// if the deadline forced them. A reader thread that later unblocks and
1545/// re-sets `closed = true` is a harmless no-op.
1546/// Timer-driven twin of [`finalize_capture_completion`] for the actor
1547/// runtime, where parking on the queue condvar would block a worker.
1548pub(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;