1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub enum ExecutableProvenance {
28 FixedNamePathSearch,
30 CallerProvided,
32}
33
34#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct DeclaredExecutable {
37 program: PathBuf,
38 provenance: ExecutableProvenance,
39}
40
41impl DeclaredExecutable {
42 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 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 pub fn program(&self) -> &Path {
74 &self.program
75 }
76
77 pub fn provenance(&self) -> ExecutableProvenance {
79 self.provenance
80 }
81
82 pub fn command(&self) -> Command {
87 Command::new(&self.program)
88 }
89}
90
91#[derive(Debug, Clone, PartialEq, Eq)]
93pub enum WorkingDirectoryPolicy {
94 Inherit,
96 Path(PathBuf),
98}
99
100#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub enum StdinPolicy {
103 Inherit,
105 Null,
107}
108
109#[derive(Debug, Clone, Copy, PartialEq, Eq)]
111pub enum OutputStreamPolicy {
112 Capture,
114 Inherit,
116}
117
118#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub enum CaptureRetention {
121 Head,
123 Tail,
125 HeadAndTail,
127}
128
129#[derive(Debug, Clone, PartialEq, Eq)]
131pub enum EnvironmentPolicy {
132 Inherit,
134 InheritWith {
136 set: BTreeMap<OsString, OsString>,
138 remove: BTreeSet<OsString>,
140 },
141 ClearAndSet {
143 set: BTreeMap<OsString, OsString>,
145 },
146}
147
148#[derive(Debug, Clone, Copy, PartialEq, Eq)]
150pub struct ProcessLimits {
151 pub timeout: Duration,
153 pub stdout_max_bytes: usize,
155 pub stderr_max_bytes: usize,
157}
158
159impl ProcessLimits {
160 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#[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
186pub const DEFAULT_DRAIN_GRACE: Duration = Duration::from_millis(250);
192
193impl ProcessSpec {
194 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 pub fn with_stdin_policy(mut self, policy: StdinPolicy) -> Self {
221 self.stdin = policy;
222 self
223 }
224
225 pub fn with_stdout_policy(mut self, policy: OutputStreamPolicy) -> Self {
227 self.stdout = policy;
228 self
229 }
230
231 pub fn with_stderr_policy(mut self, policy: OutputStreamPolicy) -> Self {
233 self.stderr = policy;
234 self
235 }
236
237 pub fn with_stdout_retention(mut self, retention: CaptureRetention) -> Self {
239 self.stdout_retention = retention;
240 self
241 }
242
243 pub fn with_stderr_retention(mut self, retention: CaptureRetention) -> Self {
245 self.stderr_retention = retention;
246 self
247 }
248
249 pub fn with_drain_grace(mut self, drain_grace: Duration) -> Self {
251 self.drain_grace = drain_grace;
252 self
253 }
254
255 pub fn executable(&self) -> &DeclaredExecutable {
257 &self.executable
258 }
259
260 pub fn arguments(&self) -> &[OsString] {
262 &self.arguments
263 }
264
265 pub fn working_directory(&self) -> &WorkingDirectoryPolicy {
267 &self.working_directory
268 }
269
270 pub fn environment(&self) -> &EnvironmentPolicy {
272 &self.environment
273 }
274
275 pub fn limits(&self) -> ProcessLimits {
277 self.limits
278 }
279
280 pub fn stdin_policy(&self) -> StdinPolicy {
282 self.stdin
283 }
284
285 pub fn stdout_policy(&self) -> OutputStreamPolicy {
287 self.stdout
288 }
289
290 pub fn stderr_policy(&self) -> OutputStreamPolicy {
292 self.stderr
293 }
294
295 pub fn stdout_retention(&self) -> CaptureRetention {
297 self.stdout_retention
298 }
299
300 pub fn stderr_retention(&self) -> CaptureRetention {
302 self.stderr_retention
303 }
304
305 pub fn drain_grace(&self) -> Duration {
307 self.drain_grace
308 }
309}
310
311#[derive(Debug, Clone, PartialEq, Eq)]
313pub struct CapturedStream {
314 pub captured: bool,
316 pub bytes: Vec<u8>,
318 pub truncated: bool,
320 pub incomplete: bool,
322 pub total_bytes: u64,
324 pub retention: CaptureRetention,
326}
327
328impl CapturedStream {
329 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#[derive(Debug)]
349pub struct ProcessOutcome {
350 pub program: PathBuf,
352 pub provenance: ExecutableProvenance,
354 pub status: ExitStatus,
356 pub timed_out: bool,
358 pub cancelled: bool,
360 pub duration: Duration,
362 pub stdout: CapturedStream,
364 pub stderr: CapturedStream,
366}
367
368impl ProcessOutcome {
369 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#[derive(Debug, Error, Clone, PartialEq, Eq)]
388#[error("{stream} was not captured completely ({reason})")]
389pub struct IncompleteOutputError {
390 pub stream: ProcessStream,
392 pub reason: IncompleteOutputReason,
394}
395
396#[derive(Debug, Clone, Copy, PartialEq, Eq)]
398pub enum IncompleteOutputReason {
399 NotCaptured,
401 Truncated,
403 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
441pub enum ProcessStream {
442 Stdout,
444 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#[derive(Debug, Error)]
459pub enum ProcessRunError {
460 #[error("process launch for {program} was cancelled before spawn")]
462 CancelledBeforeSpawn {
463 program: PathBuf,
465 },
466 #[error("failed to spawn {program} ({provenance:?}): {source}")]
468 Spawn {
469 program: PathBuf,
471 provenance: ExecutableProvenance,
473 #[source]
475 source: io::Error,
476 },
477 #[error("{stream} pipe was unavailable for {program}")]
479 MissingPipe {
480 program: PathBuf,
482 stream: ProcessStream,
484 },
485 #[error("failed to configure {stream} capture for {program}: {source}")]
487 ConfigurePipe {
488 program: PathBuf,
490 stream: ProcessStream,
492 #[source]
494 source: io::Error,
495 },
496 #[error("failed to poll {program}: {source}")]
498 Poll {
499 program: PathBuf,
501 #[source]
503 source: io::Error,
504 },
505 #[error("failed to terminate process scope for {program}: {source}")]
507 Terminate {
508 program: PathBuf,
510 #[source]
512 source: io::Error,
513 },
514 #[error("failed to reap {program}: {source}")]
516 Reap {
517 program: PathBuf,
519 #[source]
521 source: io::Error,
522 },
523 #[error("failed to drain {stream} for {program}: {source}")]
525 Capture {
526 program: PathBuf,
528 stream: ProcessStream,
530 #[source]
532 source: io::Error,
533 },
534 #[error("cleanup deadline elapsed before {program} could be reaped")]
536 CleanupDeadline {
537 program: PathBuf,
539 },
540 #[error("{feature} is unsupported on this platform")]
542 UnsupportedPlatformGuarantee {
543 feature: &'static str,
545 },
546}
547
548impl ProcessRunError {
549 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#[derive(Debug, Default, Clone, Copy)]
574pub struct ProcessRunner;
575
576#[derive(Debug, Clone, Default)]
578pub struct ProcessCancellation {
579 requested: Arc<AtomicBool>,
580}
581
582impl ProcessCancellation {
583 pub fn cancel(&self) {
585 self.requested.store(true, Ordering::SeqCst);
586 }
587
588 pub fn is_cancelled(&self) -> bool {
590 self.requested.load(Ordering::SeqCst)
591 }
592}
593
594static GLOBAL_CANCELLATION: OnceLock<ProcessCancellation> = OnceLock::new();
595
596pub fn global_cancellation_token() -> ProcessCancellation {
601 GLOBAL_CANCELLATION
602 .get_or_init(ProcessCancellation::default)
603 .clone()
604}
605
606impl ProcessRunner {
607 pub fn run(spec: &ProcessSpec) -> Result<ProcessOutcome, ProcessRunError> {
610 Self::run_cancellable(spec, &global_cancellation_token())
611 }
612
613 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 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 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 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 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 let mut info = std::mem::MaybeUninit::<libc::siginfo_t>::zeroed();
878 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 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 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 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 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#[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 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 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 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}