Skip to main content

nu_protocol/process/
child.rs

1use crate::{
2    ShellError, Span,
3    byte_stream::convert_file,
4    engine::{EngineState, FrozenJob, Job},
5    shell_error::{generic::GenericError, io::IoError},
6};
7use nu_system::{ExitStatus, ForegroundChild, ForegroundWaitStatus};
8
9use os_pipe::PipeReader;
10use std::{
11    fmt::Debug,
12    io::{self, Read},
13    sync::mpsc::{self, Receiver, RecvError, TryRecvError},
14    sync::{Arc, Mutex},
15    thread,
16};
17
18/// Check the exit status of each pipeline element.
19///
20/// This is used to implement pipefail.
21#[cfg(feature = "os")]
22pub fn check_exit_status_future(
23    exit_status: Vec<Option<ExitStatusGuard>>,
24) -> Result<(), ShellError> {
25    for one_status in exit_status.into_iter().rev().flatten() {
26        check_exit_status_future_ok(one_status)?
27    }
28    Ok(())
29}
30
31fn check_exit_status_future_ok(exit_status_guard: ExitStatusGuard) -> Result<(), ShellError> {
32    let ignore_error = {
33        let guard = exit_status_guard
34            .ignore_error
35            .lock()
36            .expect("lock ignore_error should success");
37        *guard
38    };
39    let mut future = exit_status_guard
40        .exit_status_future
41        .lock()
42        .expect("lock exit_status_future should success");
43    let span = exit_status_guard.span.unwrap_or_default();
44    let exit_status = future.wait(span)?;
45    check_ok(exit_status, ignore_error, span)
46}
47
48pub fn check_ok(status: ExitStatus, ignore_error: bool, span: Span) -> Result<(), ShellError> {
49    match status {
50        ExitStatus::Exited(exit_code) => {
51            if ignore_error {
52                Ok(())
53            } else if let Ok(exit_code) = exit_code.try_into() {
54                Err(ShellError::NonZeroExitCode { exit_code, span })
55            } else {
56                Ok(())
57            }
58        }
59        #[cfg(unix)]
60        ExitStatus::Signaled {
61            signal,
62            core_dumped,
63        } => {
64            use nix::sys::signal::Signal;
65
66            let sig = Signal::try_from(signal);
67
68            if sig == Ok(Signal::SIGPIPE) || (ignore_error && !core_dumped) {
69                // Processes often exit with SIGPIPE, but this is not an error condition.
70                Ok(())
71            } else {
72                let signal_name = sig.map(Signal::as_str).unwrap_or("unknown signal").into();
73                Err(if core_dumped {
74                    ShellError::CoreDumped {
75                        signal_name,
76                        signal,
77                        span,
78                    }
79                } else {
80                    ShellError::TerminatedBySignal {
81                        signal_name,
82                        signal,
83                        span,
84                    }
85                })
86            }
87        }
88    }
89}
90
91/// A wrapper for both `exit_status_future: Arc<Mutex<ExitStatusFuture>>`
92/// and `ignore_error: Arc<Mutex<bool>>`
93///
94/// It's useful for `pipefail` feature, which tracks exit status code with potentially
95/// ignore the error.
96#[derive(Debug)]
97pub struct ExitStatusGuard {
98    pub exit_status_future: Arc<Mutex<ExitStatusFuture>>,
99    pub ignore_error: Arc<Mutex<bool>>,
100    pub span: Option<Span>,
101}
102
103impl ExitStatusGuard {
104    pub fn new(
105        exit_status_future: Arc<Mutex<ExitStatusFuture>>,
106        ignore_error: Arc<Mutex<bool>>,
107    ) -> Self {
108        Self {
109            exit_status_future,
110            ignore_error,
111            span: None,
112        }
113    }
114
115    pub fn with_span(self, span: Span) -> Self {
116        Self {
117            exit_status_future: self.exit_status_future,
118            ignore_error: self.ignore_error,
119            span: Some(span),
120        }
121    }
122}
123
124#[derive(Debug)]
125pub enum ExitStatusFuture {
126    Finished(Result<ExitStatus, Box<ShellError>>),
127    Running(Receiver<io::Result<ExitStatus>>),
128}
129
130impl ExitStatusFuture {
131    pub fn wait(&mut self, span: Span) -> Result<ExitStatus, ShellError> {
132        match self {
133            ExitStatusFuture::Finished(Ok(status)) => Ok(*status),
134            ExitStatusFuture::Finished(Err(err)) => Err(err.as_ref().clone()),
135            ExitStatusFuture::Running(receiver) => {
136                let code = match receiver.recv() {
137                    #[cfg(unix)]
138                    Ok(Ok(
139                        status @ ExitStatus::Signaled {
140                            core_dumped: true, ..
141                        },
142                    )) => {
143                        check_ok(status, false, span)?;
144                        Ok(status)
145                    }
146                    Ok(Ok(status)) => Ok(status),
147                    Ok(Err(err)) => Err(ShellError::Io(IoError::new_with_additional_context(
148                        err,
149                        span,
150                        None,
151                        "failed to get exit code",
152                    ))),
153                    Err(err @ RecvError) => Err(ShellError::Generic(GenericError::new(
154                        err.to_string(),
155                        "failed to get exit code",
156                        span,
157                    ))),
158                };
159
160                *self = ExitStatusFuture::Finished(code.clone().map_err(Box::new));
161
162                code
163            }
164        }
165    }
166
167    fn try_wait(&mut self, span: Span) -> Result<Option<ExitStatus>, ShellError> {
168        match self {
169            ExitStatusFuture::Finished(Ok(code)) => Ok(Some(*code)),
170            ExitStatusFuture::Finished(Err(err)) => Err(err.as_ref().clone()),
171            ExitStatusFuture::Running(receiver) => {
172                let code = match receiver.try_recv() {
173                    Ok(Ok(status)) => Ok(Some(status)),
174                    Ok(Err(err)) => Err(ShellError::Generic(GenericError::new(
175                        err.to_string(),
176                        "failed to get exit code",
177                        span,
178                    ))),
179                    Err(TryRecvError::Disconnected) => Err(ShellError::Generic(GenericError::new(
180                        "receiver disconnected",
181                        "failed to get exit code",
182                        span,
183                    ))),
184                    Err(TryRecvError::Empty) => Ok(None),
185                };
186
187                if let Some(code) = code.clone().transpose() {
188                    *self = ExitStatusFuture::Finished(code.map_err(Box::new));
189                }
190
191                code
192            }
193        }
194    }
195}
196
197#[derive(derive_more::Debug)]
198pub enum ChildPipe {
199    #[debug("ChildPipe::Pipe")]
200    Pipe(PipeReader),
201
202    #[debug("ChildPipe::Tee")]
203    Tee(Box<dyn Read + Send + 'static>),
204}
205
206impl From<PipeReader> for ChildPipe {
207    fn from(pipe: PipeReader) -> Self {
208        Self::Pipe(pipe)
209    }
210}
211
212impl Read for ChildPipe {
213    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
214        match self {
215            ChildPipe::Pipe(pipe) => pipe.read(buf),
216            ChildPipe::Tee(tee) => tee.read(buf),
217        }
218    }
219}
220
221#[derive(Debug)]
222pub struct ChildProcess {
223    pub stdout: Option<ChildPipe>,
224    pub stderr: Option<ChildPipe>,
225    exit_status: Arc<Mutex<ExitStatusFuture>>,
226    ignore_error: Arc<Mutex<bool>>,
227    /// Pid of the process, when it was spawned by the shell. Lets a consumer
228    /// that stops reading early (an interactive viewer) kill a producer that
229    /// would otherwise keep the pipeline waiting on its exit status.
230    pid: Option<u32>,
231    span: Span,
232}
233
234/// A wrapper for a closure that runs once the shell finishes waiting on the process.
235#[derive(derive_more::Debug)]
236#[debug("<wait_callback>")]
237pub struct PostWaitCallback(pub Box<dyn FnOnce(ForegroundWaitStatus) + Send>);
238
239impl PostWaitCallback {
240    pub fn new<F>(f: F) -> Self
241    where
242        F: FnOnce(ForegroundWaitStatus) + Send + 'static,
243    {
244        PostWaitCallback(Box::new(f))
245    }
246
247    /// Creates a PostWaitCallback that creates a frozen job in the job table
248    /// if the incoming wait status indicates that the job was frozen.
249    ///
250    /// If `child_pid` is provided, the returned callback will also remove
251    /// it from the pid list of the current running job.
252    ///
253    /// The given `description` argument will be used as the description for the newly created job table entry.
254    pub fn for_job_control(
255        engine_state: &EngineState,
256        child_pid: Option<u32>,
257        description: Option<String>,
258    ) -> Self {
259        let this_job = engine_state.current_thread_job().cloned();
260        let jobs = engine_state.jobs.clone();
261        let is_interactive = engine_state.is_interactive;
262
263        PostWaitCallback::new(move |status| {
264            if let (Some(this_job), Some(child_pid)) = (this_job, child_pid) {
265                this_job.remove_pid(child_pid);
266            }
267
268            if let ForegroundWaitStatus::Frozen(unfreeze) = status {
269                let mut jobs = jobs.lock().expect("jobs lock is poisoned!");
270
271                let job_id = jobs.add_job(Job::Frozen(FrozenJob {
272                    unfreeze,
273                    description,
274                }));
275
276                if is_interactive {
277                    println!("\nJob {} is frozen", job_id.get());
278                }
279            }
280        })
281    }
282}
283
284impl ChildProcess {
285    pub fn new(
286        mut child: ForegroundChild,
287        reader: Option<PipeReader>,
288        swap: bool,
289        span: Span,
290        callback: Option<PostWaitCallback>,
291    ) -> Result<Self, ShellError> {
292        let (stdout, stderr) = match reader {
293            Some(combined) => (Some(combined), None),
294            None => {
295                let stdout = child.as_mut().stdout.take().map(convert_file);
296                let stderr = child.as_mut().stderr.take().map(convert_file);
297
298                if swap {
299                    (stderr, stdout)
300                } else {
301                    (stdout, stderr)
302                }
303            }
304        };
305        let pid = child.pid();
306
307        // Create a thread to wait for the exit status.
308        let (exit_status_sender, exit_status) = mpsc::channel();
309
310        thread::Builder::new()
311            .name("exit status waiter".into())
312            .spawn(move || {
313                let matched = match child.wait() {
314                    // there are two possible outcomes when we `wait` for a process to finish:
315                    // 1. the process finishes as usual
316                    // 2. (unix only) the process gets signaled with SIGTSTP
317                    //
318                    // in the second case, although the process may still be alive in a
319                    // cryonic state, we explicitly treat as it has finished with exit code 0
320                    // for the sake of the current pipeline
321                    Ok(wait_status) => {
322                        let next = match &wait_status {
323                            ForegroundWaitStatus::Frozen(_) => ExitStatus::Exited(0),
324                            ForegroundWaitStatus::Finished(exit_status) => *exit_status,
325                        };
326
327                        if let Some(callback) = callback {
328                            (callback.0)(wait_status);
329                        }
330
331                        Ok(next)
332                    }
333                    Err(err) => Err(err),
334                };
335
336                exit_status_sender.send(matched)
337            })
338            .map_err(|err| {
339                IoError::new_with_additional_context(
340                    err,
341                    span,
342                    None,
343                    "Could now spawn exit status waiter",
344                )
345            })?;
346
347        let mut this = Self::from_raw(stdout, stderr, Some(exit_status), span);
348        this.pid = Some(pid);
349        Ok(this)
350    }
351
352    pub fn from_raw(
353        stdout: Option<PipeReader>,
354        stderr: Option<PipeReader>,
355        exit_status: Option<Receiver<io::Result<ExitStatus>>>,
356        span: Span,
357    ) -> Self {
358        Self {
359            stdout: stdout.map(Into::into),
360            stderr: stderr.map(Into::into),
361            exit_status: Arc::new(Mutex::new(
362                exit_status
363                    .map(ExitStatusFuture::Running)
364                    .unwrap_or(ExitStatusFuture::Finished(Ok(ExitStatus::Exited(0)))),
365            )),
366            ignore_error: Arc::new(Mutex::new(false)),
367            pid: None,
368            span,
369        }
370    }
371
372    /// The process id, if this child was spawned by the shell.
373    pub fn pid(&self) -> Option<u32> {
374        self.pid
375    }
376
377    pub fn ignore_error(&mut self, ignore: bool) -> &mut Self {
378        {
379            let mut ignore_error = self.ignore_error.lock().expect("lock should success");
380            *ignore_error = ignore;
381        }
382        self
383    }
384
385    pub fn span(&self) -> Span {
386        self.span
387    }
388
389    pub fn into_bytes(self) -> Result<Vec<u8>, ShellError> {
390        if self.stderr.is_some() {
391            debug_assert!(false, "stderr should not exist");
392            return Err(ShellError::Generic(GenericError::new(
393                "internal error",
394                "stderr should not exist",
395                self.span,
396            )));
397        }
398
399        let bytes = (self.stdout)
400            .map(collect_bytes)
401            .transpose()
402            .map_err(|err| IoError::new(err, self.span, None))?
403            .unwrap_or_default();
404
405        let mut exit_status = self
406            .exit_status
407            .lock()
408            .expect("lock exit_status future should success");
409        let ignore_error = {
410            let guard = self
411                .ignore_error
412                .lock()
413                .expect("lock ignore error should success");
414            *guard
415        };
416        check_ok(exit_status.wait(self.span)?, ignore_error, self.span)?;
417
418        Ok(bytes)
419    }
420
421    pub fn wait(mut self) -> Result<(), ShellError> {
422        let from_io_error = IoError::factory(self.span, None);
423        if let Some(stdout) = self.stdout.take() {
424            let stderr = self
425                .stderr
426                .take()
427                .map(|stderr| {
428                    thread::Builder::new()
429                        .name("stderr consumer".into())
430                        .spawn(move || consume_pipe(stderr))
431                })
432                .transpose()
433                .map_err(&from_io_error)?;
434
435            let res = consume_pipe(stdout);
436
437            if let Some(handle) = stderr {
438                handle
439                    .join()
440                    .map_err(|e| match e.downcast::<io::Error>() {
441                        Ok(io) => from_io_error(*io).into(),
442                        Err(err) => ShellError::Generic(GenericError::new(
443                            "Unknown error",
444                            format!("{err:?}"),
445                            self.span,
446                        )),
447                    })?
448                    .map_err(&from_io_error)?;
449            }
450
451            res.map_err(&from_io_error)?;
452        } else if let Some(stderr) = self.stderr.take() {
453            consume_pipe(stderr).map_err(&from_io_error)?;
454        }
455        let mut exit_status = self
456            .exit_status
457            .lock()
458            .expect("lock exit_status future should success");
459        let ignore_error = {
460            let guard = self
461                .ignore_error
462                .lock()
463                .expect("lock ignore error should success");
464            *guard
465        };
466        check_ok(exit_status.wait(self.span)?, ignore_error, self.span)
467    }
468
469    pub fn try_wait(&mut self) -> Result<Option<ExitStatus>, ShellError> {
470        let mut exit_status = self
471            .exit_status
472            .lock()
473            .expect("lock exit_status future should success");
474        exit_status.try_wait(self.span)
475    }
476
477    pub fn wait_with_output(self) -> Result<ProcessOutput, ShellError> {
478        let from_io_error = IoError::factory(self.span, None);
479
480        let (stdout, stderr) = match (self.stdout, self.stderr) {
481            (None, None) => (None, None),
482            (None, Some(stderr)) => (None, Some(collect_bytes(stderr).map_err(&from_io_error)?)),
483            (Some(stdout), None) => (Some(collect_bytes(stdout).map_err(&from_io_error)?), None),
484            (Some(stdout), Some(stderr)) => {
485                let stderr = thread::Builder::new()
486                    .spawn(move || collect_bytes(stderr))
487                    .map_err(&from_io_error)?;
488
489                let stdout = collect_bytes(stdout).map_err(&from_io_error)?;
490
491                let stderr = stderr
492                    .join()
493                    .map_err(|e| match e.downcast::<io::Error>() {
494                        Ok(io) => from_io_error(*io).into(),
495                        Err(err) => ShellError::Generic(GenericError::new(
496                            "Unknown error",
497                            format!("{err:?}"),
498                            self.span,
499                        )),
500                    })?
501                    .map_err(&from_io_error)?;
502
503                (Some(stdout), Some(stderr))
504            }
505        };
506
507        let mut exit_status = self
508            .exit_status
509            .lock()
510            .expect("lock exit_status future should success");
511        let exit_status = exit_status.wait(self.span)?;
512
513        Ok(ProcessOutput {
514            stdout,
515            stderr,
516            exit_status,
517        })
518    }
519
520    pub fn clone_exit_status_future(&self) -> Arc<Mutex<ExitStatusFuture>> {
521        self.exit_status.clone()
522    }
523
524    pub fn clone_ignore_error(&self) -> Arc<Mutex<bool>> {
525        self.ignore_error.clone()
526    }
527}
528
529fn collect_bytes(pipe: ChildPipe) -> io::Result<Vec<u8>> {
530    let mut buf = Vec::new();
531    match pipe {
532        ChildPipe::Pipe(mut pipe) => pipe.read_to_end(&mut buf),
533        ChildPipe::Tee(mut tee) => tee.read_to_end(&mut buf),
534    }?;
535    Ok(buf)
536}
537
538fn consume_pipe(pipe: ChildPipe) -> io::Result<()> {
539    match pipe {
540        ChildPipe::Pipe(mut pipe) => io::copy(&mut pipe, &mut io::sink()),
541        ChildPipe::Tee(mut tee) => io::copy(&mut tee, &mut io::sink()),
542    }?;
543    Ok(())
544}
545
546pub struct ProcessOutput {
547    pub stdout: Option<Vec<u8>>,
548    pub stderr: Option<Vec<u8>>,
549    pub exit_status: ExitStatus,
550}