Skip to main content

running_process_platform_internal/
sync_spawn_group.rs

1use std::io;
2use std::process::{Command, Stdio};
3use std::sync::{Arc, Mutex};
4use std::thread;
5use std::time::{Duration, Instant};
6
7const DEFAULT_KILL_DRAIN_TIMEOUT: Duration = Duration::from_secs(2);
8
9fn kill_drain_deadline() -> Instant {
10    let timeout = crate::env_vars::KILL_DRAIN_TIMEOUT_MS.millis_or(DEFAULT_KILL_DRAIN_TIMEOUT);
11    Instant::now() + timeout
12}
13
14fn poll_until<T>(
15    deadline: Instant,
16    interval: Duration,
17    mut poll: impl FnMut() -> io::Result<Option<T>>,
18) -> io::Result<Option<T>> {
19    loop {
20        if let Some(value) = poll()? {
21            return Ok(Some(value));
22        }
23        let now = Instant::now();
24        if now >= deadline {
25            return Ok(None);
26        }
27        thread::sleep(interval.min(deadline.saturating_duration_since(now)));
28    }
29}
30
31trait UnixChild: Send {
32    fn kill(&mut self) -> io::Result<()>;
33    fn wait(&mut self) -> io::Result<i32>;
34    fn try_wait(&mut self) -> io::Result<Option<i32>>;
35}
36
37impl UnixChild for std::process::Child {
38    fn kill(&mut self) -> io::Result<()> {
39        std::process::Child::kill(self)
40    }
41
42    fn wait(&mut self) -> io::Result<i32> {
43        std::process::Child::wait(self).map(crate::platform::process::exit_code)
44    }
45
46    fn try_wait(&mut self) -> io::Result<Option<i32>> {
47        Ok(std::process::Child::try_wait(self)?.map(crate::platform::process::exit_code))
48    }
49}
50
51impl crate::platform::process::DaemonChildControl for std::process::Child {
52    fn kill(&mut self) -> io::Result<()> {
53        std::process::Child::kill(self)
54    }
55
56    fn wait(&mut self) -> io::Result<i32> {
57        std::process::Child::wait(self).map(crate::platform::process::exit_code)
58    }
59
60    fn try_wait(&mut self) -> io::Result<Option<i32>> {
61        Ok(std::process::Child::try_wait(self)?.map(crate::platform::process::exit_code))
62    }
63}
64
65pub struct SpawnedInner {
66    child: Arc<Mutex<Option<Box<dyn UnixChild>>>>,
67    pgid: i32,
68    retain_exit_identity: bool,
69}
70
71impl SpawnedInner {
72    // WNOWAIT observes exit without releasing the leader's PID. New spawn-mode
73    // handles retain that identity until group control is finished. If another
74    // reaper has consumed it, waitid fails and control must fail closed.
75    fn observe_owned_exit(&self) -> io::Result<Option<i32>> {
76        super::observe_owned_child_exit(self.pgid)
77    }
78
79    pub fn kill(&self) -> io::Result<()> {
80        // Try the child first, then the process group, to make sure
81        // any siblings spawned inside go down too.
82        let mut guard = self.child.lock().expect("child mutex poisoned");
83        if self.retain_exit_identity {
84            if guard.is_none() {
85                return Ok(());
86            }
87            self.observe_owned_exit()?;
88        }
89        if let Some(child) = guard.as_mut() {
90            let _ = child.kill();
91        }
92        drop(guard);
93        let _ = crate::platform::process::unix_signal_process_group(
94            self.pgid,
95            crate::platform::process::UnixSignalKind::Kill,
96        );
97        Ok(())
98    }
99
100    pub fn wait(&self) -> io::Result<i32> {
101        if self.retain_exit_identity {
102            loop {
103                if let Some(code) = self.try_wait()? {
104                    return Ok(code);
105                }
106                thread::sleep(Duration::from_millis(10));
107            }
108        }
109        let mut guard = self.child.lock().expect("child mutex poisoned");
110        let Some(child) = guard.as_mut() else {
111            return Err(io::Error::other("child handle absent"));
112        };
113        child.wait()
114    }
115
116    pub fn try_wait(&self) -> io::Result<Option<i32>> {
117        let mut guard = self.child.lock().expect("child mutex poisoned");
118        if self.retain_exit_identity {
119            return self.observe_owned_exit();
120        }
121        let Some(child) = guard.as_mut() else {
122            return Ok(None);
123        };
124        child.try_wait()
125    }
126
127    pub fn shutdown(&mut self) {
128        self.shutdown_with_deadline(kill_drain_deadline());
129    }
130
131    fn shutdown_with_deadline(&mut self, deadline: Instant) {
132        let identity_owned = !self.retain_exit_identity || self.observe_owned_exit().is_ok();
133        let group_signaled = identity_owned && crate::platform::process::unix_signal_process_group(
134            self.pgid,
135            crate::platform::process::UnixSignalKind::Kill,
136        )
137        .is_ok();
138        let Some(mut child) = self.child.lock().expect("child mutex poisoned").take() else {
139            return;
140        };
141        if !group_signaled && identity_owned {
142            let _ = child.kill();
143        }
144        match poll_until(deadline, Duration::from_millis(10), || child.try_wait()) {
145            Ok(Some(_)) => {}
146            Ok(None) | Err(_) => spawn_background_reaper(child),
147        }
148    }
149}
150
151impl crate::platform::process::SpawnedChildControl for SpawnedInner {
152    #[cfg(feature = "independent-spawn")]
153    fn retain_exit_identity(&mut self) {
154        self.retain_exit_identity = true;
155    }
156    fn kill(&mut self) -> io::Result<()> {
157        SpawnedInner::kill(self)
158    }
159
160    fn wait(&mut self) -> io::Result<i32> {
161        SpawnedInner::wait(self)
162    }
163
164    fn try_wait(&mut self) -> io::Result<Option<i32>> {
165        SpawnedInner::try_wait(self)
166    }
167
168    fn shutdown(&mut self) {
169        SpawnedInner::shutdown(self);
170    }
171}
172
173impl Drop for SpawnedInner {
174    fn drop(&mut self) {
175        // Detached handles do not invoke shutdown. Release their wait ownership
176        // without killing the process, including a retained terminal leader.
177        if self.retain_exit_identity {
178            if let Some(mut child) = self.child.lock().expect("child mutex poisoned").take() {
179                if !matches!(child.try_wait(), Ok(Some(_))) {
180                    spawn_background_reaper(child);
181                }
182            }
183        }
184    }
185}
186
187fn spawn_background_reaper(mut child: Box<dyn UnixChild>) {
188    thread::spawn(move || {
189        // Once ownership is off the caller's teardown path, a blocking wait is
190        // the most reliable terminal policy: it reaps exactly once without a
191        // retry loop, spinning, or retaining the shared child mutex.
192        let _ = child.wait();
193    });
194}
195
196fn slot_to_stdio(slot: &crate::platform::process::StdioSource<'_>) -> io::Result<Stdio> {
197    match slot {
198        crate::platform::process::StdioSource::Null => Ok(Stdio::null()),
199        crate::platform::process::StdioSource::Parent => Ok(Stdio::inherit()),
200        crate::platform::process::StdioSource::File(file) => Ok(Stdio::from(file.try_clone()?)),
201        crate::platform::process::StdioSource::Pipe => Ok(Stdio::piped()),
202    }
203}
204
205fn daemon_slot_to_stdio(
206    slot: &crate::platform::process::DaemonStdioSource<'_>,
207) -> io::Result<Stdio> {
208    match slot {
209        crate::platform::process::DaemonStdioSource::Null => Ok(Stdio::null()),
210        crate::platform::process::DaemonStdioSource::File(file) => {
211            Ok(Stdio::from(file.try_clone()?))
212        }
213    }
214}
215
216pub fn spawn_sync_daemon(
217    command: &mut Command,
218    stdio: crate::platform::process::DaemonStdio<'_>,
219    environment: crate::platform::process::SyncEnvironment,
220    _breakaway: bool,
221) -> io::Result<crate::platform::process::DaemonChild> {
222    spawn_sync_daemon_inner(command, stdio, environment, None)
223}
224
225pub fn spawn_sync_daemon_with_inheritance(
226    command: &mut Command,
227    stdio: crate::platform::process::DaemonStdio<'_>,
228    environment: crate::platform::process::SyncEnvironment,
229    _breakaway: bool,
230    inheritance: crate::platform::process::DaemonExecInheritance,
231) -> io::Result<crate::platform::process::DaemonChild> {
232    spawn_sync_daemon_inner(command, stdio, environment, Some(inheritance))
233}
234
235fn spawn_sync_daemon_inner(
236    command: &mut Command,
237    stdio: crate::platform::process::DaemonStdio<'_>,
238    environment: crate::platform::process::SyncEnvironment,
239    inheritance: Option<crate::platform::process::DaemonExecInheritance>,
240) -> io::Result<crate::platform::process::DaemonChild> {
241    apply_environment(command, environment);
242    command
243        .stdin(Stdio::null())
244        .stdout(daemon_slot_to_stdio(&stdio.stdout)?)
245        .stderr(daemon_slot_to_stdio(&stdio.stderr)?);
246
247    match inheritance {
248        Some(inheritance) => {
249            crate::platform::process::configure_sync_daemon_command_with_inheritance(
250                command,
251                inheritance,
252            )?;
253        }
254        None => crate::platform::process::configure_sync_daemon_command(command)?,
255    }
256
257    let child = crate::platform::ape::spawn_std(command, |command| command.spawn())?;
258    let pid = child.id();
259    Ok(crate::platform::process::DaemonChild {
260        pid,
261        inner: Box::new(child),
262    })
263}
264
265pub fn spawn_sync(
266    command: &mut Command,
267    stdio: crate::platform::process::SpawnStdio<'_>,
268    environment: crate::platform::process::SyncEnvironment,
269) -> io::Result<crate::platform::process::SpawnedChild> {
270    spawn_sync_inner(command, stdio, environment, false)
271}
272
273#[cfg(feature = "independent-spawn")]
274pub(crate) fn spawn_sync_owned_daemon(command: &mut Command, stdio: crate::platform::process::SpawnStdio<'_>, environment: crate::platform::process::SyncEnvironment) -> io::Result<crate::platform::process::SpawnedChild> {
275    spawn_sync_inner(command, stdio, environment, true)
276}
277
278fn spawn_sync_inner(command: &mut Command, stdio: crate::platform::process::SpawnStdio<'_>, environment: crate::platform::process::SyncEnvironment, detached: bool) -> io::Result<crate::platform::process::SpawnedChild> {
279    apply_environment(command, environment);
280    command.stdin(slot_to_stdio(&stdio.stdin)?);
281    command.stdout(slot_to_stdio(&stdio.stdout)?);
282    command.stderr(slot_to_stdio(&stdio.stderr)?);
283
284    if detached { crate::platform::process::configure_sync_daemon_command(command)?; }
285    else { crate::platform::process::configure_sync_contained_command(command)?; }
286
287    let mut child = crate::platform::ape::spawn_std(command, |command| command.spawn())?;
288    let pid = child.id();
289    let pgid = pid as i32;
290
291    let stdin = child.stdin.take();
292    let stdout = child.stdout.take();
293    let stderr = child.stderr.take();
294
295    let child: Arc<Mutex<Option<Box<dyn UnixChild>>>> = Arc::new(Mutex::new(Some(Box::new(child))));
296
297    // Drain watcher: wait for exit, then sleep `drain_timeout`. We
298    // don't proactively close anything on Unix — Rust's ChildStdin/etc.
299    // own their fds; once the child exits and the kernel ref-counts
300    // its copies to zero, parent reads will EOF naturally.
301    if let Some(timeout) = stdio.drain_timeout {
302        let child_clone = Arc::clone(&child);
303        thread::spawn(move || {
304            // Borrow child for try_wait.  We do a polling loop so
305            // shutdown() taking the inner Child during Drop doesn't
306            // wedge us.
307            loop {
308                {
309                    let mut guard = child_clone.lock().expect("child mutex poisoned");
310                    match guard.as_mut() {
311                        Some(c) => match c.try_wait() {
312                            Ok(Some(_)) => break,
313                            Ok(None) => {}
314                            Err(_) => break,
315                        },
316                        None => return,
317                    }
318                }
319                // #199: intentional — try_wait poll on the contained
320                // child, 50ms cadence inside a bounded outer drain
321                // loop. waitpid(WNOHANG)-equivalent semantics.
322                thread::sleep(std::time::Duration::from_millis(50));
323            }
324            // #199: intentional — post-mortem pipe drain. Children's
325            // write-ends of the captured stdio pipes are still being
326            // closed by the kernel after exit; this gives readers a
327            // chance to see the final bytes before the watcher
328            // releases its keep-alive.
329            thread::sleep(timeout);
330        });
331    }
332
333    Ok(crate::platform::process::SpawnedChild {
334        kill_on_drop: true,
335        stdin,
336        stdout,
337        stderr,
338        pid,
339        inner: Box::new(SpawnedInner { child, pgid, retain_exit_identity: false }),
340    })
341}
342
343fn apply_environment(
344    command: &mut Command,
345    environment: crate::platform::process::SyncEnvironment,
346) {
347    let crate::platform::process::SyncEnvironment::Explicit(base) = environment else {
348        return;
349    };
350
351    // `env_clear` also clears Command's mutation map. Preserve additions,
352    // overrides, and removals so they are replayed after the selected base.
353    let explicit: Vec<_> = command
354        .get_envs()
355        .map(|(key, value)| (key.to_os_string(), value.map(std::ffi::OsStr::to_os_string)))
356        .collect();
357    command.env_clear();
358    command.envs(base);
359    for (key, value) in explicit {
360        match value {
361            Some(value) => {
362                command.env(key, value);
363            }
364            None => {
365                command.env_remove(key);
366            }
367        }
368    }
369}
370
371#[cfg(test)]
372mod tests {
373    use super::*;
374    use std::sync::atomic::{AtomicUsize, Ordering};
375    use std::sync::{mpsc, Condvar};
376
377    struct FakeChild {
378        wait_gate: Arc<(Mutex<bool>, Condvar)>,
379        waits: Arc<AtomicUsize>,
380        kills: Arc<AtomicUsize>,
381    }
382
383    impl UnixChild for FakeChild {
384        fn kill(&mut self) -> io::Result<()> {
385            self.kills.fetch_add(1, Ordering::SeqCst);
386            Ok(())
387        }
388
389        fn wait(&mut self) -> io::Result<i32> {
390            self.waits.fetch_add(1, Ordering::SeqCst);
391            let (lock, condvar) = &*self.wait_gate;
392            let mut released = lock.lock().expect("wait gate mutex poisoned");
393            while !*released {
394                released = condvar.wait(released).expect("wait gate mutex poisoned");
395            }
396            Ok(0)
397        }
398
399        fn try_wait(&mut self) -> io::Result<Option<i32>> {
400            self.waits.fetch_add(1, Ordering::SeqCst);
401            let released = *self.wait_gate.0.lock().expect("wait gate mutex poisoned");
402            Ok(released.then_some(0))
403        }
404    }
405
406    struct BlockedFixture {
407        inner: SpawnedInner,
408        child: Arc<Mutex<Option<Box<dyn UnixChild>>>>,
409        wait_gate: Arc<(Mutex<bool>, Condvar)>,
410        waits: Arc<AtomicUsize>,
411        kills: Arc<AtomicUsize>,
412    }
413
414    fn blocked_inner() -> BlockedFixture {
415        let wait_gate = Arc::new((Mutex::new(false), Condvar::new()));
416        let waits = Arc::new(AtomicUsize::new(0));
417        let kills = Arc::new(AtomicUsize::new(0));
418        let child: Arc<Mutex<Option<Box<dyn UnixChild>>>> =
419            Arc::new(Mutex::new(Some(Box::new(FakeChild {
420                wait_gate: Arc::clone(&wait_gate),
421                waits: Arc::clone(&waits),
422                kills: Arc::clone(&kills),
423            }))));
424        BlockedFixture {
425            inner: SpawnedInner {
426                child: Arc::clone(&child),
427                pgid: i32::MAX,
428                retain_exit_identity: false,
429            },
430            child,
431            wait_gate,
432            waits,
433            kills,
434        }
435    }
436
437    fn release_wait(wait_gate: &Arc<(Mutex<bool>, Condvar)>) {
438        let (lock, condvar) = &**wait_gate;
439        *lock.lock().expect("wait gate mutex poisoned") = true;
440        condvar.notify_all();
441    }
442
443    struct ShutdownOnDrop {
444        inner: Option<SpawnedInner>,
445        deadline: Instant,
446    }
447
448    impl Drop for ShutdownOnDrop {
449        fn drop(&mut self) {
450            self.inner
451                .as_mut()
452                .expect("test wrapper missing inner")
453                .shutdown_with_deadline(self.deadline);
454        }
455    }
456
457    #[test]
458    fn drop_is_bounded_when_child_wait_does_not_complete() {
459        // Regression for #619: SpawnedChild::drop delegates directly to
460        // SpawnedInner::shutdown, modeled by this wrapper around the fake child.
461        let BlockedFixture {
462            inner, wait_gate, ..
463        } = blocked_inner();
464        let (tx, rx) = mpsc::channel();
465        let started = Instant::now();
466        let worker = thread::spawn(move || {
467            drop(ShutdownOnDrop {
468                inner: Some(inner),
469                deadline: Instant::now() + Duration::from_millis(50),
470            });
471            let _ = tx.send(started.elapsed());
472        });
473
474        // The property is causal, not a stopwatch reading: Drop must return
475        // WITHOUT waiting for the child, so it must report back before the
476        // gate is released. A blocked Drop cannot, whatever the machine load.
477        //
478        // Asserting a wall-clock bound instead conflated that with "finished
479        // inside 100ms", which a loaded runner broke by 0.2ms. The window
480        // below is generous because it only bounds how long a *failure* takes
481        // to detect: a correct Drop returns at its 50ms deadline and never
482        // approaches it.
483        let timely = rx.recv_timeout(Duration::from_secs(5));
484        release_wait(&wait_gate);
485        let returned_before_release = timely.is_ok();
486        let elapsed = timely
487            .or_else(|_| rx.recv_timeout(Duration::from_secs(5)))
488            .expect("shutdown did not unblock even after releasing fake child");
489        worker.join().expect("shutdown worker panicked");
490        assert!(
491            returned_before_release,
492            "Drop blocked in child.wait() until the fake child was released              (took {elapsed:?}); its deadline should have bounded it"
493        );
494    }
495
496    #[test]
497    fn shutdown_does_not_hold_child_mutex_while_reaping() {
498        let BlockedFixture {
499            mut inner,
500            child,
501            wait_gate,
502            waits,
503            ..
504        } = blocked_inner();
505        let worker = thread::spawn(move || {
506            inner.shutdown_with_deadline(Instant::now() + Duration::from_millis(50));
507        });
508        let deadline = Instant::now() + Duration::from_secs(1);
509        while waits.load(Ordering::SeqCst) == 0 && Instant::now() < deadline {
510            thread::yield_now();
511        }
512        // `waits` counts every try_wait poll as well as the final blocking
513        // wait: shutdown polls every 10ms until its 50ms deadline, then hands
514        // the child to a background reaper. A test thread descheduled for a
515        // poll interval or two observes more than one call, so the property
516        // is "reaping has started", not "exactly one call so far".
517        assert!(
518            waits.load(Ordering::SeqCst) >= 1,
519            "fake wait never started"
520        );
521
522        let child_mutex_available = child.try_lock().is_ok();
523        release_wait(&wait_gate);
524        worker.join().expect("shutdown worker panicked");
525        assert!(
526            child_mutex_available,
527            "shutdown held the child mutex across reaping"
528        );
529    }
530
531    struct ReadyChild {
532        polls: Arc<AtomicUsize>,
533        waits: Arc<AtomicUsize>,
534    }
535
536    impl UnixChild for ReadyChild {
537        fn kill(&mut self) -> io::Result<()> {
538            Ok(())
539        }
540
541        fn wait(&mut self) -> io::Result<i32> {
542            self.waits.fetch_add(1, Ordering::SeqCst);
543            Ok(0)
544        }
545
546        fn try_wait(&mut self) -> io::Result<Option<i32>> {
547            self.polls.fetch_add(1, Ordering::SeqCst);
548            Ok(Some(0))
549        }
550    }
551
552    #[test]
553    fn shutdown_reaps_ready_child_exactly_once() {
554        let polls = Arc::new(AtomicUsize::new(0));
555        let waits = Arc::new(AtomicUsize::new(0));
556        let child: Arc<Mutex<Option<Box<dyn UnixChild>>>> =
557            Arc::new(Mutex::new(Some(Box::new(ReadyChild {
558                polls: Arc::clone(&polls),
559                waits: Arc::clone(&waits),
560            }))));
561        let mut inner = SpawnedInner {
562            child,
563            pgid: i32::MAX,
564            retain_exit_identity: false,
565        };
566
567        inner.shutdown_with_deadline(Instant::now() + Duration::from_secs(1));
568
569        assert_eq!(polls.load(Ordering::SeqCst), 1);
570        assert_eq!(waits.load(Ordering::SeqCst), 0);
571    }
572
573    #[test]
574    fn shutdown_falls_back_to_direct_kill_when_group_signal_fails() {
575        let BlockedFixture {
576            mut inner,
577            wait_gate,
578            kills,
579            ..
580        } = blocked_inner();
581        release_wait(&wait_gate);
582
583        inner.shutdown_with_deadline(Instant::now() + Duration::from_secs(1));
584
585        assert_eq!(kills.load(Ordering::SeqCst), 1);
586    }
587}