Skip to main content

mobench_process/
lib.rs

1//! Declared executable provenance and bounded child-process policy.
2
3use std::collections::{BTreeMap, BTreeSet};
4use std::ffi::OsString;
5use std::io;
6use std::path::{Component, Path, PathBuf};
7use std::process::{Command, ExitStatus};
8use std::sync::atomic::{AtomicBool, Ordering};
9use std::sync::{Arc, OnceLock};
10use std::time::Duration;
11
12#[cfg(unix)]
13use std::collections::VecDeque;
14#[cfg(unix)]
15use std::io::Read;
16#[cfg(unix)]
17use std::os::fd::AsRawFd;
18#[cfg(unix)]
19use std::process::Stdio;
20#[cfg(unix)]
21use std::time::Instant;
22
23use thiserror::Error;
24
25/// Provenance attached to a child-process executable.
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub enum ExecutableProvenance {
28    /// A single-component executable name resolved through `PATH`.
29    FixedNamePathSearch,
30    /// An executable name or path supplied explicitly by the caller.
31    CallerProvided,
32}
33
34/// An executable whose selection source is explicit in the type.
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct DeclaredExecutable {
37    program: PathBuf,
38    provenance: ExecutableProvenance,
39}
40
41impl DeclaredExecutable {
42    /// Declare a fixed program name that Mobench may resolve through `PATH`.
43    pub fn path_search(program: impl AsRef<Path>) -> Result<Self, ExecutablePolicyError> {
44        let program = program.as_ref();
45        let mut components = program.components();
46        if program.as_os_str().is_empty()
47            || !matches!(components.next(), Some(Component::Normal(_)))
48            || components.next().is_some()
49        {
50            return Err(ExecutablePolicyError::InvalidPathSearchName(
51                program.to_path_buf(),
52            ));
53        }
54        Ok(Self {
55            program: program.to_path_buf(),
56            provenance: ExecutableProvenance::FixedNamePathSearch,
57        })
58    }
59
60    /// Accept a caller-supplied program name or path through an explicit seam.
61    pub fn explicit_override(program: impl AsRef<Path>) -> Result<Self, ExecutablePolicyError> {
62        let program = program.as_ref();
63        if program.as_os_str().is_empty() {
64            return Err(ExecutablePolicyError::EmptyExplicitOverride);
65        }
66        Ok(Self {
67            program: program.to_path_buf(),
68            provenance: ExecutableProvenance::CallerProvided,
69        })
70    }
71
72    /// Return the selected program path or fixed name.
73    pub fn program(&self) -> &Path {
74        &self.program
75    }
76
77    /// Return how this executable was selected.
78    pub fn provenance(&self) -> ExecutableProvenance {
79        self.provenance
80    }
81
82    /// Construct a raw command for the declared executable.
83    ///
84    /// Prefer [`ProcessRunner`] for production subprocesses that need bounded
85    /// output, a deadline, and explicit cwd/environment policies.
86    pub fn command(&self) -> Command {
87        Command::new(&self.program)
88    }
89}
90
91/// Explicit policy for the child process working directory.
92#[derive(Debug, Clone, PartialEq, Eq)]
93pub enum WorkingDirectoryPolicy {
94    /// Inherit the parent process working directory.
95    Inherit,
96    /// Run from the specified directory.
97    Path(PathBuf),
98}
99
100/// Explicit policy for the child process standard input stream.
101#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub enum StdinPolicy {
103    /// Connect the child to the parent's standard input.
104    Inherit,
105    /// Connect the child to the platform null device.
106    Null,
107}
108
109/// Explicit policy for one child output stream.
110#[derive(Debug, Clone, Copy, PartialEq, Eq)]
111pub enum OutputStreamPolicy {
112    /// Capture the stream with the configured byte and retention bounds.
113    Capture,
114    /// Connect the child directly to the corresponding parent stream.
115    Inherit,
116}
117
118/// Which portions of an oversized captured stream are retained.
119#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub enum CaptureRetention {
121    /// Retain the beginning of the stream.
122    Head,
123    /// Retain the end of the stream.
124    Tail,
125    /// Split the byte budget between the beginning and end of the stream.
126    HeadAndTail,
127}
128
129/// Explicit policy for the child process environment.
130#[derive(Debug, Clone, PartialEq, Eq)]
131pub enum EnvironmentPolicy {
132    /// Inherit the parent environment unchanged.
133    Inherit,
134    /// Inherit the parent environment, then remove and set selected variables.
135    InheritWith {
136        /// Variables to set after applying removals.
137        set: BTreeMap<OsString, OsString>,
138        /// Variables to remove from the inherited environment.
139        remove: BTreeSet<OsString>,
140    },
141    /// Clear the parent environment and set only the supplied variables.
142    ClearAndSet {
143        /// Complete environment visible to the child.
144        set: BTreeMap<OsString, OsString>,
145    },
146}
147
148/// Output and wall-clock bounds for one child process.
149#[derive(Debug, Clone, Copy, PartialEq, Eq)]
150pub struct ProcessLimits {
151    /// Maximum wall-clock time before the child is killed and reaped.
152    pub timeout: Duration,
153    /// Maximum stdout bytes retained in memory.
154    pub stdout_max_bytes: usize,
155    /// Maximum stderr bytes retained in memory.
156    pub stderr_max_bytes: usize,
157}
158
159impl ProcessLimits {
160    /// Construct explicit process bounds.
161    pub const fn new(timeout: Duration, stdout_max_bytes: usize, stderr_max_bytes: usize) -> Self {
162        Self {
163            timeout,
164            stdout_max_bytes,
165            stderr_max_bytes,
166        }
167    }
168}
169
170/// Complete specification for one bounded child-process invocation.
171#[derive(Debug, Clone, PartialEq, Eq)]
172pub struct ProcessSpec {
173    executable: DeclaredExecutable,
174    arguments: Vec<OsString>,
175    working_directory: WorkingDirectoryPolicy,
176    environment: EnvironmentPolicy,
177    limits: ProcessLimits,
178    stdin: StdinPolicy,
179    stdout: OutputStreamPolicy,
180    stderr: OutputStreamPolicy,
181    stdout_retention: CaptureRetention,
182    stderr_retention: CaptureRetention,
183    drain_grace: Duration,
184}
185
186/// Default time allowed to collect pipe data after process-scope termination.
187///
188/// Once this grace elapses, capture pipes are closed and the returned stream is
189/// marked incomplete. This keeps escaped descendants from extending the
190/// process deadline indefinitely.
191pub const DEFAULT_DRAIN_GRACE: Duration = Duration::from_millis(250);
192
193impl ProcessSpec {
194    /// Construct a process specification with explicit cwd, environment, and
195    /// resource limits. Standard input defaults to null; stdout and stderr
196    /// default to bounded head-and-tail capture for source compatibility.
197    pub fn new(
198        executable: DeclaredExecutable,
199        arguments: Vec<OsString>,
200        working_directory: WorkingDirectoryPolicy,
201        environment: EnvironmentPolicy,
202        limits: ProcessLimits,
203    ) -> Self {
204        Self {
205            executable,
206            arguments,
207            working_directory,
208            environment,
209            limits,
210            stdin: StdinPolicy::Null,
211            stdout: OutputStreamPolicy::Capture,
212            stderr: OutputStreamPolicy::Capture,
213            stdout_retention: CaptureRetention::HeadAndTail,
214            stderr_retention: CaptureRetention::HeadAndTail,
215            drain_grace: DEFAULT_DRAIN_GRACE,
216        }
217    }
218
219    /// Select the child standard-input policy.
220    pub fn with_stdin_policy(mut self, policy: StdinPolicy) -> Self {
221        self.stdin = policy;
222        self
223    }
224
225    /// Select the child standard-output policy.
226    pub fn with_stdout_policy(mut self, policy: OutputStreamPolicy) -> Self {
227        self.stdout = policy;
228        self
229    }
230
231    /// Select the child standard-error policy.
232    pub fn with_stderr_policy(mut self, policy: OutputStreamPolicy) -> Self {
233        self.stderr = policy;
234        self
235    }
236
237    /// Select retention for captured standard output.
238    pub fn with_stdout_retention(mut self, retention: CaptureRetention) -> Self {
239        self.stdout_retention = retention;
240        self
241    }
242
243    /// Select retention for captured standard error.
244    pub fn with_stderr_retention(mut self, retention: CaptureRetention) -> Self {
245        self.stderr_retention = retention;
246        self
247    }
248
249    /// Override the bounded post-termination pipe-drain grace.
250    pub fn with_drain_grace(mut self, drain_grace: Duration) -> Self {
251        self.drain_grace = drain_grace;
252        self
253    }
254
255    /// Return the declared executable.
256    pub fn executable(&self) -> &DeclaredExecutable {
257        &self.executable
258    }
259
260    /// Return the literal argv entries passed to the executable.
261    pub fn arguments(&self) -> &[OsString] {
262        &self.arguments
263    }
264
265    /// Return the working-directory policy.
266    pub fn working_directory(&self) -> &WorkingDirectoryPolicy {
267        &self.working_directory
268    }
269
270    /// Return the environment policy.
271    pub fn environment(&self) -> &EnvironmentPolicy {
272        &self.environment
273    }
274
275    /// Return the process resource limits.
276    pub fn limits(&self) -> ProcessLimits {
277        self.limits
278    }
279
280    /// Return the child standard-input policy.
281    pub fn stdin_policy(&self) -> StdinPolicy {
282        self.stdin
283    }
284
285    /// Return the child standard-output policy.
286    pub fn stdout_policy(&self) -> OutputStreamPolicy {
287        self.stdout
288    }
289
290    /// Return the child standard-error policy.
291    pub fn stderr_policy(&self) -> OutputStreamPolicy {
292        self.stderr
293    }
294
295    /// Return standard-output capture retention.
296    pub fn stdout_retention(&self) -> CaptureRetention {
297        self.stdout_retention
298    }
299
300    /// Return standard-error capture retention.
301    pub fn stderr_retention(&self) -> CaptureRetention {
302        self.stderr_retention
303    }
304
305    /// Return the bounded post-termination pipe-drain grace.
306    pub fn drain_grace(&self) -> Duration {
307        self.drain_grace
308    }
309}
310
311/// Bounded output captured from one child stream.
312#[derive(Debug, Clone, PartialEq, Eq)]
313pub struct CapturedStream {
314    /// Whether this stream was configured for capture.
315    pub captured: bool,
316    /// Retained bytes according to `retention`, never larger than the cap.
317    pub bytes: Vec<u8>,
318    /// Whether bytes were omitted from the retained capture.
319    pub truncated: bool,
320    /// Whether capture ended before EOF because the drain grace elapsed.
321    pub incomplete: bool,
322    /// Bytes observed before EOF or forced pipe closure.
323    pub total_bytes: u64,
324    /// Retention policy used for this stream.
325    pub retention: CaptureRetention,
326}
327
328impl CapturedStream {
329    /// True when the bytes cannot be treated as the complete child stream.
330    pub const fn is_partial(&self) -> bool {
331        self.truncated || self.incomplete
332    }
333
334    #[cfg(unix)]
335    fn inherited(retention: CaptureRetention) -> Self {
336        Self {
337            captured: false,
338            bytes: Vec::new(),
339            truncated: false,
340            incomplete: false,
341            total_bytes: 0,
342            retention,
343        }
344    }
345}
346
347/// Structured result from a child that was successfully spawned and reaped.
348#[derive(Debug)]
349pub struct ProcessOutcome {
350    /// Program name/path selected by the caller.
351    pub program: PathBuf,
352    /// Provenance of the executable selection.
353    pub provenance: ExecutableProvenance,
354    /// Reaped child exit status, including the post-kill status on timeout.
355    pub status: ExitStatus,
356    /// True when the deadline elapsed and the runner killed the child.
357    pub timed_out: bool,
358    /// True when cooperative cancellation killed the child process scope.
359    pub cancelled: bool,
360    /// Wall-clock duration including output drain and child reap.
361    pub duration: Duration,
362    /// Bounded stdout capture.
363    pub stdout: CapturedStream,
364    /// Bounded stderr capture.
365    pub stderr: CapturedStream,
366}
367
368impl ProcessOutcome {
369    /// Convert a captured outcome into the standard library's complete output
370    /// shape.
371    ///
372    /// Machine-readable callers must use this instead of parsing retained
373    /// bytes directly: inherited, truncated, or incompletely drained streams
374    /// are rejected rather than being mistaken for a complete document.
375    pub fn into_complete_output(self) -> Result<std::process::Output, IncompleteOutputError> {
376        validate_complete_stream(&self.stdout, ProcessStream::Stdout)?;
377        validate_complete_stream(&self.stderr, ProcessStream::Stderr)?;
378        Ok(std::process::Output {
379            status: self.status,
380            stdout: self.stdout.bytes,
381            stderr: self.stderr.bytes,
382        })
383    }
384}
385
386/// A captured child stream was not a complete machine-readable value.
387#[derive(Debug, Error, Clone, PartialEq, Eq)]
388#[error("{stream} was not captured completely ({reason})")]
389pub struct IncompleteOutputError {
390    /// Stream that did not contain a complete value.
391    pub stream: ProcessStream,
392    /// Why the stream cannot be parsed as complete output.
393    pub reason: IncompleteOutputReason,
394}
395
396/// Why captured output cannot be treated as complete.
397#[derive(Debug, Clone, Copy, PartialEq, Eq)]
398pub enum IncompleteOutputReason {
399    /// The stream was connected directly to the parent.
400    NotCaptured,
401    /// The configured byte budget omitted part of the stream.
402    Truncated,
403    /// The drain grace elapsed before the stream reached EOF.
404    Incomplete,
405}
406
407impl std::fmt::Display for IncompleteOutputReason {
408    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
409        match self {
410            Self::NotCaptured => formatter.write_str("stream was inherited"),
411            Self::Truncated => formatter.write_str("capture byte limit was exceeded"),
412            Self::Incomplete => formatter.write_str("capture ended before EOF"),
413        }
414    }
415}
416
417fn validate_complete_stream(
418    stream: &CapturedStream,
419    kind: ProcessStream,
420) -> Result<(), IncompleteOutputError> {
421    let reason = if !stream.captured {
422        Some(IncompleteOutputReason::NotCaptured)
423    } else if stream.incomplete {
424        Some(IncompleteOutputReason::Incomplete)
425    } else if stream.truncated {
426        Some(IncompleteOutputReason::Truncated)
427    } else {
428        None
429    };
430    match reason {
431        Some(reason) => Err(IncompleteOutputError {
432            stream: kind,
433            reason,
434        }),
435        None => Ok(()),
436    }
437}
438
439/// Identifies one captured child stream in structured errors.
440#[derive(Debug, Clone, Copy, PartialEq, Eq)]
441pub enum ProcessStream {
442    /// Child standard output.
443    Stdout,
444    /// Child standard error.
445    Stderr,
446}
447
448impl std::fmt::Display for ProcessStream {
449    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
450        match self {
451            Self::Stdout => formatter.write_str("stdout"),
452            Self::Stderr => formatter.write_str("stderr"),
453        }
454    }
455}
456
457/// Structured failure while starting, supervising, or capturing a child.
458#[derive(Debug, Error)]
459pub enum ProcessRunError {
460    /// Cooperative cancellation was already requested before spawn.
461    #[error("process launch for {program} was cancelled before spawn")]
462    CancelledBeforeSpawn {
463        /// Declared program name/path.
464        program: PathBuf,
465    },
466    /// The child could not be started.
467    #[error("failed to spawn {program} ({provenance:?}): {source}")]
468    Spawn {
469        /// Declared program name/path.
470        program: PathBuf,
471        /// Executable-selection provenance.
472        provenance: ExecutableProvenance,
473        /// Operating-system error.
474        #[source]
475        source: io::Error,
476    },
477    /// A configured pipe was unexpectedly unavailable.
478    #[error("{stream} pipe was unavailable for {program}")]
479    MissingPipe {
480        /// Declared program name/path.
481        program: PathBuf,
482        /// Missing stream.
483        stream: ProcessStream,
484    },
485    /// A capture pipe could not be switched to nonblocking mode.
486    #[error("failed to configure {stream} capture for {program}: {source}")]
487    ConfigurePipe {
488        /// Declared program name/path.
489        program: PathBuf,
490        /// Stream that could not be configured.
491        stream: ProcessStream,
492        /// Operating-system error.
493        #[source]
494        source: io::Error,
495    },
496    /// The child status could not be polled.
497    #[error("failed to poll {program}: {source}")]
498    Poll {
499        /// Declared program name/path.
500        program: PathBuf,
501        /// Operating-system error.
502        #[source]
503        source: io::Error,
504    },
505    /// The child process scope could not be terminated during cleanup.
506    #[error("failed to terminate process scope for {program}: {source}")]
507    Terminate {
508        /// Declared program name/path.
509        program: PathBuf,
510        /// Operating-system error.
511        #[source]
512        source: io::Error,
513    },
514    /// The child could not be reaped.
515    #[error("failed to reap {program}: {source}")]
516    Reap {
517        /// Declared program name/path.
518        program: PathBuf,
519        /// Operating-system error.
520        #[source]
521        source: io::Error,
522    },
523    /// Draining a child stream failed.
524    #[error("failed to drain {stream} for {program}: {source}")]
525    Capture {
526        /// Declared program name/path.
527        program: PathBuf,
528        /// Stream that failed.
529        stream: ProcessStream,
530        /// Operating-system error.
531        #[source]
532        source: io::Error,
533    },
534    /// The direct child did not become reapable within the cleanup grace.
535    #[error("cleanup deadline elapsed before {program} could be reaped")]
536    CleanupDeadline {
537        /// Declared program name/path.
538        program: PathBuf,
539    },
540    /// The platform cannot provide the process-tree deadline guarantee.
541    #[error("{feature} is unsupported on this platform")]
542    UnsupportedPlatformGuarantee {
543        /// Guarantee that is unavailable.
544        feature: &'static str,
545    },
546}
547
548impl ProcessRunError {
549    /// Return the operating-system error kind when spawning failed.
550    pub fn spawn_error_kind(&self) -> Option<io::ErrorKind> {
551        match self {
552            Self::Spawn { source, .. } => Some(source.kind()),
553            Self::CancelledBeforeSpawn { .. }
554            | Self::MissingPipe { .. }
555            | Self::ConfigurePipe { .. }
556            | Self::Poll { .. }
557            | Self::Terminate { .. }
558            | Self::Reap { .. }
559            | Self::Capture { .. }
560            | Self::CleanupDeadline { .. }
561            | Self::UnsupportedPlatformGuarantee { .. } => None,
562        }
563    }
564}
565
566/// Synchronous bounded child-process runner.
567///
568/// On Unix, each child is placed in a new process group. Capture pipes are
569/// nonblocking and are closed after a bounded drain grace, so even a descendant
570/// that escapes the group cannot hold [`ProcessRunner::run`] open indefinitely.
571/// Platforms without this implementation fail before spawning with
572/// [`ProcessRunError::UnsupportedPlatformGuarantee`].
573#[derive(Debug, Default, Clone, Copy)]
574pub struct ProcessRunner;
575
576/// Thread-safe cooperative process cancellation signal.
577#[derive(Debug, Clone, Default)]
578pub struct ProcessCancellation {
579    requested: Arc<AtomicBool>,
580}
581
582impl ProcessCancellation {
583    /// Request cancellation of current and subsequent runs using this token.
584    pub fn cancel(&self) {
585        self.requested.store(true, Ordering::SeqCst);
586    }
587
588    /// Whether cancellation has been requested.
589    pub fn is_cancelled(&self) -> bool {
590        self.requested.load(Ordering::SeqCst)
591    }
592}
593
594static GLOBAL_CANCELLATION: OnceLock<ProcessCancellation> = OnceLock::new();
595
596/// Return the process-wide token used by [`ProcessRunner::run`].
597///
598/// CLI entrypoints can connect an OS interruption handler to this token so
599/// subprocesses launched through the CLI and SDK share one cancellation scope.
600pub fn global_cancellation_token() -> ProcessCancellation {
601    GLOBAL_CANCELLATION
602        .get_or_init(ProcessCancellation::default)
603        .clone()
604}
605
606impl ProcessRunner {
607    /// Run the specified child, concurrently drain stdout/stderr, enforce the
608    /// deadline, and return only after the child has been reaped.
609    pub fn run(spec: &ProcessSpec) -> Result<ProcessOutcome, ProcessRunError> {
610        Self::run_cancellable(spec, &global_cancellation_token())
611    }
612
613    /// Run with an explicit cooperative cancellation token.
614    pub fn run_cancellable(
615        spec: &ProcessSpec,
616        cancellation: &ProcessCancellation,
617    ) -> Result<ProcessOutcome, ProcessRunError> {
618        if cancellation.is_cancelled() {
619            return Err(ProcessRunError::CancelledBeforeSpawn {
620                program: spec.executable.program().to_path_buf(),
621            });
622        }
623        run_platform(spec, cancellation)
624    }
625}
626
627#[cfg(not(unix))]
628fn run_platform(
629    _spec: &ProcessSpec,
630    _cancellation: &ProcessCancellation,
631) -> Result<ProcessOutcome, ProcessRunError> {
632    Err(ProcessRunError::UnsupportedPlatformGuarantee {
633        feature: "bounded process-tree supervision",
634    })
635}
636
637#[cfg(unix)]
638fn run_platform(
639    spec: &ProcessSpec,
640    cancellation: &ProcessCancellation,
641) -> Result<ProcessOutcome, ProcessRunError> {
642    const MAX_POLL_INTERVAL: Duration = Duration::from_millis(10);
643
644    let program = spec.executable.program().to_path_buf();
645    let provenance = spec.executable.provenance();
646    let started = Instant::now();
647    let mut command = spec.executable.command();
648    command.args(&spec.arguments);
649    apply_working_directory_policy(&mut command, &spec.working_directory);
650    apply_environment_policy(&mut command, &spec.environment);
651    configure_process_group(&mut command);
652    apply_stdio_policy(&mut command, spec);
653
654    let child = command.spawn().map_err(|source| ProcessRunError::Spawn {
655        program: program.clone(),
656        provenance,
657        source,
658    })?;
659    let mut scope = ChildScope::new(child);
660
661    let mut stdout = match spec.stdout {
662        OutputStreamPolicy::Capture => {
663            let pipe = scope
664                .child
665                .stdout
666                .take()
667                .ok_or_else(|| ProcessRunError::MissingPipe {
668                    program: program.clone(),
669                    stream: ProcessStream::Stdout,
670                })?;
671            Some(ActiveCapture::new(
672                pipe,
673                spec.limits.stdout_max_bytes,
674                spec.stdout_retention,
675                &program,
676                ProcessStream::Stdout,
677            )?)
678        }
679        OutputStreamPolicy::Inherit => None,
680    };
681    let mut stderr = match spec.stderr {
682        OutputStreamPolicy::Capture => {
683            let pipe = scope
684                .child
685                .stderr
686                .take()
687                .ok_or_else(|| ProcessRunError::MissingPipe {
688                    program: program.clone(),
689                    stream: ProcessStream::Stderr,
690                })?;
691            Some(ActiveCapture::new(
692                pipe,
693                spec.limits.stderr_max_bytes,
694                spec.stderr_retention,
695                &program,
696                ProcessStream::Stderr,
697            )?)
698        }
699        OutputStreamPolicy::Inherit => None,
700    };
701
702    let mut status = None;
703    let mut timed_out = false;
704    let mut cancelled = false;
705    let mut cleanup_started = None;
706
707    loop {
708        drain_optional_capture(&mut stdout, &program, ProcessStream::Stdout)?;
709        drain_optional_capture(&mut stderr, &program, ProcessStream::Stderr)?;
710
711        if status.is_none()
712            && scope.has_exited().map_err(|source| ProcessRunError::Poll {
713                program: program.clone(),
714                source,
715            })?
716        {
717            // Keep the direct child waitable until the process group has been
718            // terminated. The unreaped child pins its PID/PGID, preventing a
719            // successful short-lived command from racing an unrelated process
720            // group that later reuses the numeric id.
721            scope
722                .terminate(true)
723                .map_err(|source| ProcessRunError::Terminate {
724                    program: program.clone(),
725                    source,
726                })?;
727            status = Some(scope.reap().map_err(|source| ProcessRunError::Reap {
728                program: program.clone(),
729                source,
730            })?);
731        }
732
733        if status.is_some() && cleanup_started.is_none() {
734            cleanup_started = Some(Instant::now());
735        } else if status.is_none() && cleanup_started.is_none() && cancellation.is_cancelled() {
736            cancelled = true;
737            scope
738                .terminate(false)
739                .map_err(|source| ProcessRunError::Terminate {
740                    program: program.clone(),
741                    source,
742                })?;
743            cleanup_started = Some(Instant::now());
744        } else if status.is_none()
745            && cleanup_started.is_none()
746            && started.elapsed() >= spec.limits.timeout
747        {
748            timed_out = true;
749            scope
750                .terminate(false)
751                .map_err(|source| ProcessRunError::Terminate {
752                    program: program.clone(),
753                    source,
754                })?;
755            cleanup_started = Some(Instant::now());
756        }
757
758        let captures_complete = capture_complete(&stdout) && capture_complete(&stderr);
759        if let Some(exit_status) = status
760            && captures_complete
761        {
762            return Ok(build_outcome(
763                &program,
764                OutcomeState {
765                    provenance,
766                    status: exit_status,
767                    timed_out,
768                    cancelled,
769                    started,
770                },
771                (stdout, stderr),
772                spec,
773                false,
774            ));
775        }
776
777        if let Some(cleanup_started) = cleanup_started
778            && cleanup_started.elapsed() >= spec.drain_grace
779        {
780            let Some(exit_status) = status else {
781                return Err(ProcessRunError::CleanupDeadline { program });
782            };
783            return Ok(build_outcome(
784                &program,
785                OutcomeState {
786                    provenance,
787                    status: exit_status,
788                    timed_out,
789                    cancelled,
790                    started,
791                },
792                (stdout, stderr),
793                spec,
794                true,
795            ));
796        }
797
798        let wait = next_poll_interval(started, cleanup_started, spec, MAX_POLL_INTERVAL);
799        wait_for_pipe_activity(&stdout, &stderr, wait, &program)?;
800    }
801}
802
803#[cfg(unix)]
804fn configure_process_group(command: &mut Command) {
805    use std::os::unix::process::CommandExt;
806
807    // A zero process-group id makes the child the leader of a fresh group.
808    command.process_group(0);
809}
810
811#[cfg(unix)]
812#[derive(Debug, Clone, Copy)]
813struct ProcessGroup {
814    id: libc::pid_t,
815}
816
817#[cfg(unix)]
818impl ProcessGroup {
819    fn for_child(child: &std::process::Child) -> Self {
820        Self {
821            id: child.id() as libc::pid_t,
822        }
823    }
824
825    fn terminate(self, child: &mut std::process::Child, direct_exited: bool) -> io::Result<()> {
826        // SAFETY: `id` is the positive pid returned by `Child::id`; negating it
827        // addresses only the process group created for this child.
828        let result = unsafe { libc::kill(-self.id, libc::SIGKILL) };
829        let group_error = (result != 0).then(io::Error::last_os_error);
830
831        if !direct_exited {
832            match child.kill() {
833                Ok(()) => {}
834                Err(error) if error.kind() == io::ErrorKind::InvalidInput => {}
835                Err(error) => return Err(error),
836            }
837        }
838
839        match group_error {
840            // Darwin reports EPERM when the process group contains only the
841            // unreaped zombie leader. If a live same-uid descendant remained
842            // in the group, it would still be signalable and killpg would
843            // succeed. The waitable leader continues to pin the numeric PGID.
844            Some(error) if direct_exited && error.kind() == io::ErrorKind::PermissionDenied => {
845                Ok(())
846            }
847            Some(error) if error.raw_os_error() != Some(libc::ESRCH) => Err(error),
848            Some(_) | None => Ok(()),
849        }
850    }
851}
852
853#[cfg(unix)]
854struct ChildScope {
855    child: std::process::Child,
856    process_group: ProcessGroup,
857    terminated: bool,
858    status: Option<ExitStatus>,
859}
860
861#[cfg(unix)]
862impl ChildScope {
863    fn new(child: std::process::Child) -> Self {
864        let process_group = ProcessGroup::for_child(&child);
865        Self {
866            child,
867            process_group,
868            terminated: false,
869            status: None,
870        }
871    }
872
873    fn has_exited(&self) -> io::Result<bool> {
874        // `waitid(..., WNOWAIT)` observes the terminal state without reaping
875        // the child. Holding the zombie until the process group is terminated
876        // pins the numeric PID/PGID and removes the reuse race around killpg.
877        let mut info = std::mem::MaybeUninit::<libc::siginfo_t>::zeroed();
878        // SAFETY: `info` points to writable siginfo storage and the positive
879        // pid belongs to the live `Child` handle. WNOWAIT preserves waitability.
880        let result = unsafe {
881            libc::waitid(
882                libc::P_PID,
883                self.process_group.id as libc::id_t,
884                info.as_mut_ptr(),
885                libc::WEXITED | libc::WNOHANG | libc::WNOWAIT,
886            )
887        };
888        if result != 0 {
889            let error = io::Error::last_os_error();
890            if error.kind() == io::ErrorKind::Interrupted {
891                return Ok(false);
892            }
893            return Err(error);
894        }
895        // SAFETY: waitid initialized `info` on success, including the no-state
896        // case where si_pid is zero.
897        Ok(unsafe { info.assume_init().si_pid() } == self.process_group.id)
898    }
899
900    fn reap(&mut self) -> io::Result<ExitStatus> {
901        if let Some(status) = self.status {
902            return Ok(status);
903        }
904        let status = self.child.wait()?;
905        self.status = Some(status);
906        Ok(status)
907    }
908
909    fn terminate(&mut self, direct_exited: bool) -> io::Result<()> {
910        if !self.terminated {
911            self.process_group
912                .terminate(&mut self.child, direct_exited)?;
913            self.terminated = true;
914        }
915        Ok(())
916    }
917}
918
919#[cfg(unix)]
920impl Drop for ChildScope {
921    fn drop(&mut self) {
922        if !self.terminated {
923            let _ = self
924                .process_group
925                .terminate(&mut self.child, self.status.is_some());
926        }
927        if self.status.is_none() {
928            let _ = self.child.kill();
929            // SIGKILL is not catchable. A blocking wait here is therefore
930            // bounded by kernel task teardown and guarantees that error paths
931            // do not leak a zombie direct child.
932            let _ = self.child.wait();
933        }
934    }
935}
936
937#[cfg(unix)]
938fn apply_working_directory_policy(command: &mut Command, policy: &WorkingDirectoryPolicy) {
939    match policy {
940        WorkingDirectoryPolicy::Inherit => {}
941        WorkingDirectoryPolicy::Path(path) => {
942            command.current_dir(path);
943        }
944    }
945}
946
947#[cfg(unix)]
948fn apply_environment_policy(command: &mut Command, policy: &EnvironmentPolicy) {
949    match policy {
950        EnvironmentPolicy::Inherit => {}
951        EnvironmentPolicy::InheritWith { set, remove } => {
952            for key in remove {
953                command.env_remove(key);
954            }
955            command.envs(set);
956        }
957        EnvironmentPolicy::ClearAndSet { set } => {
958            command.env_clear().envs(set);
959        }
960    }
961}
962
963#[cfg(unix)]
964fn apply_stdio_policy(command: &mut Command, spec: &ProcessSpec) {
965    command.stdin(match spec.stdin {
966        StdinPolicy::Inherit => Stdio::inherit(),
967        StdinPolicy::Null => Stdio::null(),
968    });
969    command.stdout(match spec.stdout {
970        OutputStreamPolicy::Capture => Stdio::piped(),
971        OutputStreamPolicy::Inherit => Stdio::inherit(),
972    });
973    command.stderr(match spec.stderr {
974        OutputStreamPolicy::Capture => Stdio::piped(),
975        OutputStreamPolicy::Inherit => Stdio::inherit(),
976    });
977}
978
979#[cfg(unix)]
980struct ActiveCapture<R> {
981    reader: R,
982    accumulator: CaptureAccumulator,
983    eof: bool,
984}
985
986#[cfg(unix)]
987impl<R: Read + AsRawFd> ActiveCapture<R> {
988    fn new(
989        reader: R,
990        max_bytes: usize,
991        retention: CaptureRetention,
992        program: &Path,
993        stream: ProcessStream,
994    ) -> Result<Self, ProcessRunError> {
995        set_nonblocking(reader.as_raw_fd()).map_err(|source| ProcessRunError::ConfigurePipe {
996            program: program.to_path_buf(),
997            stream,
998            source,
999        })?;
1000        Ok(Self {
1001            reader,
1002            accumulator: CaptureAccumulator::new(max_bytes, retention),
1003            eof: false,
1004        })
1005    }
1006
1007    fn drain_available(&mut self) -> io::Result<()> {
1008        const DRAIN_QUANTUM_BYTES: usize = 256 * 1024;
1009
1010        let mut buffer = [0_u8; 8 * 1024];
1011        let mut drained = 0_usize;
1012        while drained < DRAIN_QUANTUM_BYTES {
1013            match self.reader.read(&mut buffer) {
1014                Ok(0) => {
1015                    self.eof = true;
1016                    return Ok(());
1017                }
1018                Ok(read) => {
1019                    self.accumulator.record(&buffer[..read]);
1020                    drained = drained.saturating_add(read);
1021                }
1022                Err(error) if error.kind() == io::ErrorKind::WouldBlock => return Ok(()),
1023                Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1024                Err(error) => return Err(error),
1025            }
1026        }
1027        Ok(())
1028    }
1029
1030    fn finish(self, force_incomplete: bool) -> CapturedStream {
1031        self.accumulator.finish(force_incomplete && !self.eof)
1032    }
1033}
1034
1035#[cfg(unix)]
1036fn drain_optional_capture<R: Read + AsRawFd>(
1037    capture: &mut Option<ActiveCapture<R>>,
1038    program: &Path,
1039    stream: ProcessStream,
1040) -> Result<(), ProcessRunError> {
1041    if let Some(capture) = capture {
1042        capture
1043            .drain_available()
1044            .map_err(|source| ProcessRunError::Capture {
1045                program: program.to_path_buf(),
1046                stream,
1047                source,
1048            })?;
1049    }
1050    Ok(())
1051}
1052
1053#[cfg(unix)]
1054fn capture_complete<R>(capture: &Option<ActiveCapture<R>>) -> bool {
1055    capture.as_ref().is_none_or(|capture| capture.eof)
1056}
1057
1058#[cfg(unix)]
1059struct OutcomeState {
1060    provenance: ExecutableProvenance,
1061    status: ExitStatus,
1062    timed_out: bool,
1063    cancelled: bool,
1064    started: Instant,
1065}
1066
1067#[cfg(unix)]
1068fn build_outcome(
1069    program: &Path,
1070    state: OutcomeState,
1071    captures: (
1072        Option<ActiveCapture<std::process::ChildStdout>>,
1073        Option<ActiveCapture<std::process::ChildStderr>>,
1074    ),
1075    spec: &ProcessSpec,
1076    force_incomplete: bool,
1077) -> ProcessOutcome {
1078    let (stdout, stderr) = captures;
1079    ProcessOutcome {
1080        program: program.to_path_buf(),
1081        provenance: state.provenance,
1082        status: state.status,
1083        timed_out: state.timed_out,
1084        cancelled: state.cancelled,
1085        duration: state.started.elapsed(),
1086        stdout: stdout.map_or_else(
1087            || CapturedStream::inherited(spec.stdout_retention),
1088            |capture| capture.finish(force_incomplete),
1089        ),
1090        stderr: stderr.map_or_else(
1091            || CapturedStream::inherited(spec.stderr_retention),
1092            |capture| capture.finish(force_incomplete),
1093        ),
1094    }
1095}
1096
1097#[cfg(unix)]
1098fn next_poll_interval(
1099    started: Instant,
1100    cleanup_started: Option<Instant>,
1101    spec: &ProcessSpec,
1102    maximum: Duration,
1103) -> Duration {
1104    match cleanup_started {
1105        Some(started) => maximum.min(spec.drain_grace.saturating_sub(started.elapsed())),
1106        None => maximum.min(spec.limits.timeout.saturating_sub(started.elapsed())),
1107    }
1108}
1109
1110#[cfg(unix)]
1111fn wait_for_pipe_activity<Stdout: AsRawFd, Stderr: AsRawFd>(
1112    stdout: &Option<ActiveCapture<Stdout>>,
1113    stderr: &Option<ActiveCapture<Stderr>>,
1114    timeout: Duration,
1115    program: &Path,
1116) -> Result<(), ProcessRunError> {
1117    let mut descriptors = Vec::with_capacity(2);
1118    if let Some(capture) = stdout
1119        && !capture.eof
1120    {
1121        descriptors.push(libc::pollfd {
1122            fd: capture.reader.as_raw_fd(),
1123            events: libc::POLLIN | libc::POLLHUP | libc::POLLERR,
1124            revents: 0,
1125        });
1126    }
1127    if let Some(capture) = stderr
1128        && !capture.eof
1129    {
1130        descriptors.push(libc::pollfd {
1131            fd: capture.reader.as_raw_fd(),
1132            events: libc::POLLIN | libc::POLLHUP | libc::POLLERR,
1133            revents: 0,
1134        });
1135    }
1136
1137    let timeout_millis = timeout
1138        .as_nanos()
1139        .saturating_add(999_999)
1140        .saturating_div(1_000_000)
1141        .min(i32::MAX as u128) as libc::c_int;
1142    if descriptors.is_empty() {
1143        std::thread::sleep(timeout);
1144        return Ok(());
1145    }
1146
1147    // SAFETY: `descriptors` is a valid mutable pollfd array for the duration of
1148    // this call, and its length is passed unchanged.
1149    let result = unsafe {
1150        libc::poll(
1151            descriptors.as_mut_ptr(),
1152            descriptors.len() as libc::nfds_t,
1153            timeout_millis,
1154        )
1155    };
1156    if result >= 0 {
1157        return Ok(());
1158    }
1159    let error = io::Error::last_os_error();
1160    if error.kind() == io::ErrorKind::Interrupted {
1161        Ok(())
1162    } else {
1163        Err(ProcessRunError::Poll {
1164            program: program.to_path_buf(),
1165            source: error,
1166        })
1167    }
1168}
1169
1170#[cfg(unix)]
1171fn set_nonblocking(fd: std::os::fd::RawFd) -> io::Result<()> {
1172    // SAFETY: `fd` is owned by a live child pipe. `F_GETFL` does not mutate
1173    // memory and `F_SETFL` updates only flags on that descriptor.
1174    let flags = unsafe { libc::fcntl(fd, libc::F_GETFL) };
1175    if flags < 0 {
1176        return Err(io::Error::last_os_error());
1177    }
1178    let result = unsafe { libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) };
1179    if result < 0 {
1180        Err(io::Error::last_os_error())
1181    } else {
1182        Ok(())
1183    }
1184}
1185
1186#[cfg(unix)]
1187struct CaptureAccumulator {
1188    head: Vec<u8>,
1189    tail: VecDeque<u8>,
1190    max_bytes: usize,
1191    total_bytes: u64,
1192    retention: CaptureRetention,
1193}
1194
1195#[cfg(unix)]
1196impl CaptureAccumulator {
1197    fn new(max_bytes: usize, retention: CaptureRetention) -> Self {
1198        Self {
1199            head: Vec::with_capacity(max_bytes.min(8 * 1024)),
1200            tail: VecDeque::new(),
1201            max_bytes,
1202            total_bytes: 0,
1203            retention,
1204        }
1205    }
1206
1207    fn record(&mut self, bytes: &[u8]) {
1208        self.total_bytes = self.total_bytes.saturating_add(bytes.len() as u64);
1209        match self.retention {
1210            CaptureRetention::Head => {
1211                let retained = self
1212                    .max_bytes
1213                    .saturating_sub(self.head.len())
1214                    .min(bytes.len());
1215                self.head.extend_from_slice(&bytes[..retained]);
1216            }
1217            CaptureRetention::Tail => retain_tail(&mut self.tail, bytes, self.max_bytes),
1218            CaptureRetention::HeadAndTail => {
1219                let head_max = self.max_bytes.saturating_add(1) / 2;
1220                let head_retained = head_max.saturating_sub(self.head.len()).min(bytes.len());
1221                self.head.extend_from_slice(&bytes[..head_retained]);
1222                retain_tail(
1223                    &mut self.tail,
1224                    &bytes[head_retained..],
1225                    self.max_bytes.saturating_sub(head_max),
1226                );
1227            }
1228        }
1229    }
1230
1231    fn finish(self, incomplete: bool) -> CapturedStream {
1232        let mut bytes = self.head;
1233        bytes.extend(self.tail);
1234        CapturedStream {
1235            captured: true,
1236            truncated: incomplete || self.total_bytes > bytes.len() as u64,
1237            incomplete,
1238            total_bytes: self.total_bytes,
1239            retention: self.retention,
1240            bytes,
1241        }
1242    }
1243}
1244
1245#[cfg(unix)]
1246fn retain_tail(tail: &mut VecDeque<u8>, bytes: &[u8], maximum: usize) {
1247    if maximum == 0 {
1248        return;
1249    }
1250    if bytes.len() >= maximum {
1251        tail.clear();
1252        tail.extend(&bytes[bytes.len() - maximum..]);
1253        return;
1254    }
1255    let overflow = tail
1256        .len()
1257        .saturating_add(bytes.len())
1258        .saturating_sub(maximum);
1259    tail.drain(..overflow);
1260    tail.extend(bytes);
1261}
1262
1263/// Executable-selection policy failure.
1264#[derive(Debug, Error, PartialEq, Eq)]
1265pub enum ExecutablePolicyError {
1266    #[error("executable must be one fixed PATH-search name, got {0}")]
1267    InvalidPathSearchName(PathBuf),
1268    #[error("explicit executable override cannot be empty")]
1269    EmptyExplicitOverride,
1270}
1271
1272#[cfg(test)]
1273mod tests {
1274    use super::*;
1275
1276    #[cfg(unix)]
1277    use std::fs;
1278    #[cfg(unix)]
1279    use std::os::unix::fs::PermissionsExt;
1280
1281    #[cfg(unix)]
1282    fn write_executable_script(directory: &Path, name: &str, body: &str) -> PathBuf {
1283        let path = directory.join(name);
1284        // Publish the executable atomically after writing and chmod'ing it. On
1285        // Linux, executing a path while another test process is still mutating
1286        // that same inode can transiently return ETXTBSY ("Text file busy").
1287        // Keeping the writable staging inode private avoids making unrelated
1288        // process-runner tests depend on filesystem timing.
1289        let staging_path = directory.join(format!(".{name}.staging"));
1290        fs::write(&staging_path, format!("#!/bin/sh\nset -eu\n{body}\n"))
1291            .expect("write test script");
1292        let mut permissions = fs::metadata(&staging_path)
1293            .expect("script metadata")
1294            .permissions();
1295        permissions.set_mode(0o755);
1296        fs::set_permissions(&staging_path, permissions).expect("make test script executable");
1297        fs::rename(&staging_path, &path).expect("publish test script");
1298        path
1299    }
1300
1301    #[cfg(unix)]
1302    fn test_spec(
1303        executable: DeclaredExecutable,
1304        arguments: Vec<OsString>,
1305        timeout: Duration,
1306        stdout_max_bytes: usize,
1307        stderr_max_bytes: usize,
1308    ) -> ProcessSpec {
1309        ProcessSpec::new(
1310            executable,
1311            arguments,
1312            WorkingDirectoryPolicy::Inherit,
1313            EnvironmentPolicy::Inherit,
1314            ProcessLimits::new(timeout, stdout_max_bytes, stderr_max_bytes),
1315        )
1316    }
1317
1318    #[test]
1319    fn built_in_path_search_accepts_one_fixed_name() {
1320        let executable = DeclaredExecutable::path_search("python3").expect("declare executable");
1321
1322        assert_eq!(executable.program(), Path::new("python3"));
1323        assert_eq!(
1324            executable.provenance(),
1325            ExecutableProvenance::FixedNamePathSearch
1326        );
1327    }
1328
1329    #[test]
1330    fn built_in_path_search_rejects_paths_and_empty_names() {
1331        for invalid in ["", "./python", "tools/python", "../python", "/bin/python"] {
1332            assert!(matches!(
1333                DeclaredExecutable::path_search(invalid),
1334                Err(ExecutablePolicyError::InvalidPathSearchName(_))
1335            ));
1336        }
1337    }
1338
1339    #[test]
1340    fn explicit_override_records_caller_provenance() {
1341        let executable =
1342            DeclaredExecutable::explicit_override("./tools/python").expect("declare override");
1343
1344        assert_eq!(executable.program(), Path::new("./tools/python"));
1345        assert_eq!(
1346            executable.provenance(),
1347            ExecutableProvenance::CallerProvided
1348        );
1349    }
1350
1351    #[cfg(unix)]
1352    #[test]
1353    fn runner_passes_metacharacters_as_literal_argv() {
1354        let root = tempfile::tempdir().expect("tempdir");
1355        let script = write_executable_script(root.path(), "literal argv", "printf '%s' \"$1\"");
1356        let marker = root.path().join("shell-expanded");
1357        let payload = format!("literal;$(touch {}) `false` * ?", marker.display());
1358        let spec = test_spec(
1359            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1360            vec![payload.clone().into()],
1361            Duration::from_secs(5),
1362            4 * 1024,
1363            4 * 1024,
1364        );
1365
1366        let outcome = ProcessRunner::run(&spec).expect("run literal argv script");
1367
1368        assert!(outcome.status.success());
1369        assert_eq!(outcome.stdout.bytes, payload.as_bytes());
1370        assert!(!marker.exists(), "argv was interpreted by a shell");
1371    }
1372
1373    #[cfg(unix)]
1374    #[test]
1375    fn runner_applies_working_directory_and_clear_environment() {
1376        let root = tempfile::tempdir().expect("tempdir");
1377        let script = write_executable_script(
1378            root.path(),
1379            "cwd-env",
1380            "printf '%s|%s|%s' \"$MOBENCH_VALUE\" \"${HOME-unset}\" \"$(pwd)\"",
1381        );
1382        let mut set = BTreeMap::new();
1383        set.insert(OsString::from("MOBENCH_VALUE"), OsString::from("visible"));
1384        let spec = ProcessSpec::new(
1385            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1386            Vec::new(),
1387            WorkingDirectoryPolicy::Path(root.path().to_path_buf()),
1388            EnvironmentPolicy::ClearAndSet { set },
1389            ProcessLimits::new(Duration::from_secs(5), 4096, 4096),
1390        );
1391
1392        let outcome = ProcessRunner::run(&spec).expect("run with cwd and clear environment");
1393        let stdout = String::from_utf8_lossy(&outcome.stdout.bytes);
1394        let mut fields = stdout.split('|');
1395        assert_eq!(fields.next(), Some("visible"));
1396        assert_eq!(fields.next(), Some("unset"));
1397        let reported_cwd = PathBuf::from(fields.next().expect("reported cwd"));
1398        assert_eq!(
1399            reported_cwd.canonicalize().expect("canonical reported cwd"),
1400            root.path().canonicalize().expect("canonical expected cwd")
1401        );
1402        assert!(!outcome.stdout.is_partial());
1403    }
1404
1405    #[cfg(unix)]
1406    #[test]
1407    fn runner_applies_inherit_with_set_and_remove() {
1408        let root = tempfile::tempdir().expect("tempdir");
1409        let script = write_executable_script(
1410            root.path(),
1411            "inherit-with",
1412            "printf '%s|%s' \"$MOBENCH_PROCESS_VISIBLE\" \"${HOME-unset}\"",
1413        );
1414        let mut set = BTreeMap::new();
1415        set.insert(
1416            OsString::from("MOBENCH_PROCESS_VISIBLE"),
1417            OsString::from("yes"),
1418        );
1419        let mut remove = BTreeSet::new();
1420        remove.insert(OsString::from("HOME"));
1421        let spec = ProcessSpec::new(
1422            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1423            Vec::new(),
1424            WorkingDirectoryPolicy::Inherit,
1425            EnvironmentPolicy::InheritWith { set, remove },
1426            ProcessLimits::new(Duration::from_secs(5), 4096, 4096),
1427        );
1428
1429        let outcome = ProcessRunner::run(&spec).expect("run with edited environment");
1430        assert_eq!(outcome.stdout.bytes, b"yes|unset");
1431    }
1432
1433    #[cfg(unix)]
1434    #[test]
1435    fn runner_uses_null_stdin_and_supports_inherited_output() {
1436        let root = tempfile::tempdir().expect("tempdir");
1437        let script = write_executable_script(
1438            root.path(),
1439            "stdin-null",
1440            "if IFS= read -r line; then exit 9; else printf eof >&2; fi",
1441        );
1442        let spec = test_spec(
1443            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1444            Vec::new(),
1445            Duration::from_secs(5),
1446            1024,
1447            1024,
1448        )
1449        .with_stdin_policy(StdinPolicy::Null)
1450        .with_stdout_policy(OutputStreamPolicy::Inherit);
1451
1452        let outcome = ProcessRunner::run(&spec).expect("run with null stdin");
1453        assert!(outcome.status.success());
1454        assert!(!outcome.stdout.captured);
1455        assert_eq!(outcome.stderr.bytes, b"eof");
1456    }
1457
1458    #[cfg(unix)]
1459    #[test]
1460    fn runner_zero_caps_drain_and_report_truncation() {
1461        let root = tempfile::tempdir().expect("tempdir");
1462        let script =
1463            write_executable_script(root.path(), "zero-caps", "printf stdout; printf stderr >&2");
1464        let spec = test_spec(
1465            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1466            Vec::new(),
1467            Duration::from_secs(5),
1468            0,
1469            0,
1470        );
1471
1472        let outcome = ProcessRunner::run(&spec).expect("drain zero-cap streams");
1473        assert!(outcome.stdout.bytes.is_empty());
1474        assert!(outcome.stderr.bytes.is_empty());
1475        assert_eq!(outcome.stdout.total_bytes, 6);
1476        assert_eq!(outcome.stderr.total_bytes, 6);
1477        assert!(outcome.stdout.truncated);
1478        assert!(outcome.stderr.truncated);
1479        assert!(!outcome.stdout.incomplete);
1480        assert!(!outcome.stderr.incomplete);
1481    }
1482
1483    #[cfg(unix)]
1484    #[test]
1485    fn runner_retains_head_and_tail_for_machine_detectable_truncation() {
1486        let root = tempfile::tempdir().expect("tempdir");
1487        let script = write_executable_script(root.path(), "head-tail", "printf 0123456789");
1488        let spec = test_spec(
1489            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1490            Vec::new(),
1491            Duration::from_secs(5),
1492            6,
1493            0,
1494        );
1495
1496        let outcome = ProcessRunner::run(&spec).expect("capture head and tail");
1497        assert_eq!(outcome.stdout.bytes, b"012789");
1498        assert_eq!(outcome.stdout.total_bytes, 10);
1499        assert_eq!(outcome.stdout.retention, CaptureRetention::HeadAndTail);
1500        assert!(outcome.stdout.is_partial());
1501        assert!(!outcome.stdout.incomplete);
1502    }
1503
1504    #[cfg(unix)]
1505    #[test]
1506    fn runner_caps_dual_streams_without_deadlock() {
1507        let root = tempfile::tempdir().expect("tempdir");
1508        let script = write_executable_script(
1509            root.path(),
1510            "dual-stream",
1511            "dd if=/dev/zero bs=4096 count=64 2>/dev/null\n\
1512             dd if=/dev/zero bs=4096 count=64 1>&2 2>/dev/null",
1513        );
1514        let spec = test_spec(
1515            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1516            Vec::new(),
1517            Duration::from_secs(5),
1518            1024,
1519            2048,
1520        );
1521
1522        let outcome = ProcessRunner::run(&spec).expect("drain both streams");
1523
1524        assert!(outcome.status.success());
1525        assert_eq!(outcome.stdout.bytes.len(), 1024);
1526        assert_eq!(outcome.stderr.bytes.len(), 2048);
1527        assert_eq!(outcome.stdout.total_bytes, 64 * 4096);
1528        assert_eq!(outcome.stderr.total_bytes, 64 * 4096);
1529        assert!(outcome.stdout.truncated);
1530        assert!(outcome.stderr.truncated);
1531    }
1532
1533    #[cfg(unix)]
1534    #[test]
1535    fn runner_reports_truncation_per_stream() {
1536        let root = tempfile::tempdir().expect("tempdir");
1537        let script = write_executable_script(
1538            root.path(),
1539            "one-truncated-stream",
1540            "printf 'short'\n\
1541             dd if=/dev/zero bs=4096 count=2 1>&2 2>/dev/null",
1542        );
1543        let spec = test_spec(
1544            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1545            Vec::new(),
1546            Duration::from_secs(5),
1547            1024,
1548            512,
1549        );
1550
1551        let outcome = ProcessRunner::run(&spec).expect("capture streams");
1552
1553        assert_eq!(outcome.stdout.bytes, b"short");
1554        assert!(!outcome.stdout.truncated);
1555        assert_eq!(outcome.stderr.bytes.len(), 512);
1556        assert!(outcome.stderr.truncated);
1557    }
1558
1559    #[cfg(unix)]
1560    #[test]
1561    fn complete_machine_output_rejects_truncated_and_inherited_streams() {
1562        let root = tempfile::tempdir().expect("tempdir");
1563        let script = write_executable_script(root.path(), "machine-output", "printf 0123456789");
1564        let truncated = test_spec(
1565            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1566            Vec::new(),
1567            Duration::from_secs(5),
1568            4,
1569            1024,
1570        );
1571
1572        let error = ProcessRunner::run(&truncated)
1573            .expect("run truncated output script")
1574            .into_complete_output()
1575            .expect_err("truncated machine output must be rejected");
1576        assert_eq!(error.stream, ProcessStream::Stdout);
1577        assert_eq!(error.reason, IncompleteOutputReason::Truncated);
1578
1579        let inherited = test_spec(
1580            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1581            Vec::new(),
1582            Duration::from_secs(5),
1583            1024,
1584            1024,
1585        )
1586        .with_stdout_policy(OutputStreamPolicy::Inherit);
1587        let error = ProcessRunner::run(&inherited)
1588            .expect("run inherited output script")
1589            .into_complete_output()
1590            .expect_err("inherited machine output must be rejected");
1591        assert_eq!(error.stream, ProcessStream::Stdout);
1592        assert_eq!(error.reason, IncompleteOutputReason::NotCaptured);
1593    }
1594
1595    #[cfg(unix)]
1596    #[test]
1597    fn complete_machine_output_preserves_status_and_bytes() {
1598        let root = tempfile::tempdir().expect("tempdir");
1599        let script = write_executable_script(
1600            root.path(),
1601            "complete-output",
1602            "printf stdout; printf stderr >&2; exit 7",
1603        );
1604        let spec = test_spec(
1605            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1606            Vec::new(),
1607            Duration::from_secs(5),
1608            1024,
1609            1024,
1610        );
1611
1612        let output = ProcessRunner::run(&spec)
1613            .expect("run complete output script")
1614            .into_complete_output()
1615            .expect("convert complete output");
1616        assert_eq!(output.status.code(), Some(7));
1617        assert_eq!(output.stdout, b"stdout");
1618        assert_eq!(output.stderr, b"stderr");
1619    }
1620
1621    #[cfg(unix)]
1622    #[test]
1623    fn runner_times_out_kills_and_reaps_child() {
1624        let root = tempfile::tempdir().expect("tempdir");
1625        let script = write_executable_script(
1626            root.path(),
1627            "hang-with-descendant",
1628            "sleep 30 &\nwhile :; do :; done",
1629        );
1630        let spec = test_spec(
1631            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1632            Vec::new(),
1633            Duration::from_millis(250),
1634            1024,
1635            1024,
1636        );
1637
1638        let outcome = ProcessRunner::run(&spec).expect("time out hanging child");
1639
1640        assert!(outcome.timed_out);
1641        assert!(!outcome.cancelled);
1642        assert!(!outcome.status.success());
1643        assert!(
1644            outcome.duration < Duration::from_secs(5),
1645            "deadline enforcement took {:?}",
1646            outcome.duration
1647        );
1648    }
1649
1650    #[cfg(unix)]
1651    #[test]
1652    fn cooperative_cancellation_kills_and_reaps_the_process_scope() {
1653        let root = tempfile::tempdir().expect("tempdir");
1654        let script = write_executable_script(
1655            root.path(),
1656            "cancel-with-descendant",
1657            "sleep 30 &\nwhile :; do :; done",
1658        );
1659        let spec = test_spec(
1660            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1661            Vec::new(),
1662            Duration::from_secs(30),
1663            1024,
1664            1024,
1665        );
1666        let cancellation = ProcessCancellation::default();
1667        let request = cancellation.clone();
1668        let canceller = std::thread::spawn(move || {
1669            std::thread::sleep(Duration::from_millis(100));
1670            request.cancel();
1671        });
1672
1673        let outcome = ProcessRunner::run_cancellable(&spec, &cancellation)
1674            .expect("cooperatively cancel child");
1675        canceller.join().expect("join cancellation request");
1676
1677        assert!(outcome.cancelled);
1678        assert!(!outcome.timed_out);
1679        assert!(!outcome.status.success());
1680        assert!(outcome.duration < Duration::from_secs(2));
1681    }
1682
1683    #[test]
1684    fn cancellation_requested_before_run_fails_before_spawn() {
1685        let cancellation = ProcessCancellation::default();
1686        cancellation.cancel();
1687        let spec = ProcessSpec::new(
1688            DeclaredExecutable::path_search("definitely-not-executed").expect("declare program"),
1689            Vec::new(),
1690            WorkingDirectoryPolicy::Inherit,
1691            EnvironmentPolicy::Inherit,
1692            ProcessLimits::new(Duration::from_secs(1), 1024, 1024),
1693        );
1694
1695        assert!(matches!(
1696            ProcessRunner::run_cancellable(&spec, &cancellation),
1697            Err(ProcessRunError::CancelledBeforeSpawn { .. })
1698        ));
1699    }
1700
1701    #[cfg(unix)]
1702    #[test]
1703    fn continuous_dual_stream_output_cannot_starve_deadline() {
1704        let root = tempfile::tempdir().expect("tempdir");
1705        let script = write_executable_script(
1706            root.path(),
1707            "continuous-output",
1708            "while :; do printf 0123456789; printf 9876543210 >&2; done",
1709        );
1710        let spec = test_spec(
1711            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1712            Vec::new(),
1713            Duration::from_millis(150),
1714            64,
1715            64,
1716        );
1717
1718        let outcome = ProcessRunner::run(&spec).expect("bound continuous output");
1719        assert!(outcome.timed_out);
1720        assert!(outcome.duration < Duration::from_secs(2));
1721        assert!(outcome.stdout.bytes.len() <= 64);
1722        assert!(outcome.stderr.bytes.len() <= 64);
1723    }
1724
1725    #[cfg(unix)]
1726    #[test]
1727    fn runner_bounds_escaped_session_descendant_pipe_drain() {
1728        const MODE: &str = "MOBENCH_PROCESS_ESCAPED_DESCENDANT_MODE";
1729        const PID_PATH: &str = "MOBENCH_PROCESS_ESCAPED_DESCENDANT_PID_PATH";
1730        const READY_PATH: &str = "MOBENCH_PROCESS_ESCAPED_DESCENDANT_READY_PATH";
1731        const TEST_NAME: &str = "tests::runner_bounds_escaped_session_descendant_pipe_drain";
1732
1733        match std::env::var(MODE).as_deref() {
1734            Ok("child") => {
1735                // SAFETY: this subprocess exists only to verify that a new
1736                // session can retain inherited capture descriptors.
1737                assert_ne!(unsafe { libc::setsid() }, -1, "setsid failed");
1738                fs::write(std::env::var_os(READY_PATH).expect("ready path"), b"ready")
1739                    .expect("write ready marker");
1740                std::thread::sleep(Duration::from_secs(10));
1741                return;
1742            }
1743            Ok("parent") => {
1744                let child = Command::new(std::env::current_exe().expect("current test exe"))
1745                    .arg(TEST_NAME)
1746                    .arg("--exact")
1747                    .arg("--nocapture")
1748                    .env(MODE, "child")
1749                    .env(
1750                        READY_PATH,
1751                        std::env::var_os(READY_PATH).expect("ready path"),
1752                    )
1753                    .stdout(Stdio::inherit())
1754                    .stderr(Stdio::inherit())
1755                    .spawn()
1756                    .expect("spawn escaped descendant");
1757                fs::write(
1758                    std::env::var_os(PID_PATH).expect("pid path"),
1759                    child.id().to_string(),
1760                )
1761                .expect("write escaped pid");
1762                let ready_path = PathBuf::from(
1763                    std::env::var_os(READY_PATH).expect("ready path should remain set"),
1764                );
1765                let wait_started = Instant::now();
1766                while !ready_path.exists() && wait_started.elapsed() < Duration::from_secs(2) {
1767                    std::thread::sleep(Duration::from_millis(10));
1768                }
1769                assert!(ready_path.exists(), "escaped child did not become ready");
1770                drop(child);
1771                return;
1772            }
1773            _ => {}
1774        }
1775
1776        struct KillEscapedChild(libc::pid_t);
1777        impl Drop for KillEscapedChild {
1778            fn drop(&mut self) {
1779                // SAFETY: the pid was written by the dedicated child process.
1780                let _ = unsafe { libc::kill(self.0, libc::SIGKILL) };
1781            }
1782        }
1783
1784        let root = tempfile::tempdir().expect("tempdir");
1785        let pid_path = root.path().join("escaped.pid");
1786        let ready_path = root.path().join("escaped.ready");
1787        let mut set = BTreeMap::new();
1788        set.insert(OsString::from(MODE), OsString::from("parent"));
1789        set.insert(OsString::from(PID_PATH), pid_path.clone().into_os_string());
1790        set.insert(OsString::from(READY_PATH), ready_path.into_os_string());
1791        let spec = ProcessSpec::new(
1792            DeclaredExecutable::explicit_override(
1793                std::env::current_exe().expect("current test executable"),
1794            )
1795            .expect("declare test executable"),
1796            vec![TEST_NAME.into(), "--exact".into(), "--nocapture".into()],
1797            WorkingDirectoryPolicy::Inherit,
1798            EnvironmentPolicy::InheritWith {
1799                set,
1800                remove: BTreeSet::new(),
1801            },
1802            ProcessLimits::new(Duration::from_secs(3), 4096, 4096),
1803        )
1804        .with_drain_grace(Duration::from_millis(150));
1805
1806        let outcome = ProcessRunner::run(&spec);
1807        let escaped_pid = fs::read_to_string(&pid_path)
1808            .expect("read escaped pid")
1809            .parse::<libc::pid_t>()
1810            .expect("parse escaped pid");
1811        let _cleanup = KillEscapedChild(escaped_pid);
1812        let outcome = outcome.expect("return despite escaped pipe holder");
1813
1814        assert!(outcome.status.success());
1815        assert!(!outcome.timed_out);
1816        assert!(outcome.duration < Duration::from_secs(2));
1817        assert!(outcome.stdout.incomplete);
1818        assert!(outcome.stderr.incomplete);
1819        assert!(outcome.stdout.is_partial());
1820        assert!(outcome.stderr.is_partial());
1821    }
1822
1823    #[cfg(unix)]
1824    #[test]
1825    fn runner_cleans_pipe_holding_descendant_after_parent_exit() {
1826        let root = tempfile::tempdir().expect("tempdir");
1827        let script = write_executable_script(
1828            root.path(),
1829            "exiting-parent",
1830            "sleep 30 &\nprintf parent-exited",
1831        );
1832        let spec = test_spec(
1833            DeclaredExecutable::explicit_override(&script).expect("declare script"),
1834            Vec::new(),
1835            Duration::from_secs(2),
1836            1024,
1837            1024,
1838        );
1839
1840        let outcome = ProcessRunner::run(&spec).expect("clean up pipe-holding descendant");
1841
1842        assert!(outcome.status.success());
1843        assert!(!outcome.timed_out);
1844        assert_eq!(outcome.stdout.bytes, b"parent-exited");
1845        assert!(
1846            outcome.duration < Duration::from_secs(2),
1847            "descendant held capture pipes for {:?}",
1848            outcome.duration
1849        );
1850    }
1851
1852    #[cfg(unix)]
1853    #[test]
1854    fn runner_preserves_executable_provenance() {
1855        let path_search = test_spec(
1856            DeclaredExecutable::path_search("sh").expect("declare PATH search"),
1857            vec!["-c".into(), "printf path-search".into()],
1858            Duration::from_secs(5),
1859            1024,
1860            1024,
1861        );
1862        let explicit = test_spec(
1863            DeclaredExecutable::explicit_override("/bin/sh").expect("declare explicit shell"),
1864            vec!["-c".into(), "printf explicit".into()],
1865            Duration::from_secs(5),
1866            1024,
1867            1024,
1868        );
1869
1870        let path_outcome = ProcessRunner::run(&path_search).expect("run PATH search");
1871        let explicit_outcome = ProcessRunner::run(&explicit).expect("run explicit executable");
1872
1873        assert_eq!(
1874            path_outcome.provenance,
1875            ExecutableProvenance::FixedNamePathSearch
1876        );
1877        assert_eq!(
1878            explicit_outcome.provenance,
1879            ExecutableProvenance::CallerProvided
1880        );
1881        assert_eq!(path_outcome.stdout.bytes, b"path-search");
1882        assert_eq!(explicit_outcome.stdout.bytes, b"explicit");
1883    }
1884
1885    #[cfg(unix)]
1886    #[test]
1887    fn ambient_executable_selection_does_not_replace_declared_program() {
1888        let root = tempfile::tempdir().expect("tempdir");
1889        let marker = root.path().join("ambient-ran");
1890        let declared = write_executable_script(root.path(), "declared", "printf declared");
1891        let ambient = write_executable_script(
1892            root.path(),
1893            "ambient",
1894            &format!("printf ambient > {}", marker.display()),
1895        );
1896        let mut set = BTreeMap::new();
1897        set.insert(
1898            OsString::from("MOBENCH_PROCESS_EXECUTABLE"),
1899            ambient.into_os_string(),
1900        );
1901        let spec = ProcessSpec::new(
1902            DeclaredExecutable::explicit_override(&declared).expect("declare executable"),
1903            Vec::new(),
1904            WorkingDirectoryPolicy::Inherit,
1905            EnvironmentPolicy::InheritWith {
1906                set,
1907                remove: BTreeSet::new(),
1908            },
1909            ProcessLimits::new(Duration::from_secs(5), 1024, 1024),
1910        );
1911
1912        let outcome = ProcessRunner::run(&spec).expect("run declared executable");
1913
1914        assert_eq!(outcome.stdout.bytes, b"declared");
1915        assert!(!marker.exists(), "ambient executable was invoked");
1916    }
1917
1918    #[cfg(not(unix))]
1919    #[test]
1920    fn unsupported_platform_fails_before_spawning() {
1921        let spec = ProcessSpec::new(
1922            DeclaredExecutable::path_search("definitely-not-executed").expect("declare program"),
1923            Vec::new(),
1924            WorkingDirectoryPolicy::Inherit,
1925            EnvironmentPolicy::Inherit,
1926            ProcessLimits::new(Duration::from_secs(1), 1024, 1024),
1927        );
1928
1929        assert!(matches!(
1930            ProcessRunner::run(&spec),
1931            Err(ProcessRunError::UnsupportedPlatformGuarantee { .. })
1932        ));
1933    }
1934}