running_process_platform_internal/
sync_spawn_group.rs1use 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 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 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 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 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 if let Some(timeout) = stdio.drain_timeout {
302 let child_clone = Arc::clone(&child);
303 thread::spawn(move || {
304 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 thread::sleep(std::time::Duration::from_millis(50));
323 }
324 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 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 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 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 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}