Skip to main content

command_stream/execa/
process.rs

1use super::result::strip_final_newline;
2use super::{CancelSignal, ExecaCommand, ExecaError, ExecaOutcome, ExecaResult, Options};
3use crate::{OutputChunk, OutputStream, StreamingRunner};
4use std::future::pending;
5use std::time::Instant;
6use tokio::task::JoinHandle;
7use tokio::time::{sleep_until, Instant as TokioInstant};
8
9/// Live output and process control. Dropping the handle requests termination.
10pub struct Subprocess {
11    stream: OutputStream,
12    task: JoinHandle<crate::Result<()>>,
13    options: Options,
14    result: ExecaResult,
15    started: Instant,
16    deadline: Option<TokioInstant>,
17    stopped: bool,
18}
19
20impl Subprocess {
21    pub(crate) fn new(command: ExecaCommand) -> Self {
22        let options = command.options;
23        let display: Vec<_> = std::iter::once(&command.file)
24            .chain(&command.args)
25            .map(|arg| arg.to_string_lossy().into_owned())
26            .collect();
27        let escaped_command = display
28            .iter()
29            .map(|arg| crate::quote(arg))
30            .collect::<Vec<_>>()
31            .join(" ");
32        let mut runner = StreamingRunner::from_argv(command.file, command.args)
33            .env(options.env.clone())
34            .clear_env(!options.extend_env)
35            .prefer_local(options.prefer_local.clone())
36            .kill_signal(options.kill_signal.clone());
37        if let Some(cwd) = &options.cwd {
38            runner = runner.cwd(cwd);
39        }
40        if let Some(input) = &options.input {
41            runner = runner.stdin_bytes(input.clone());
42        }
43        let (stream, task) = runner.spawn();
44        crate::trace::trace_lazy("Execa", || "Started exact-argv command".to_string());
45        Self {
46            stream,
47            task,
48            deadline: options.timeout.map(|timeout| TokioInstant::now() + timeout),
49            result: ExecaResult {
50                command: display.join(" "),
51                escaped_command,
52                all: options.all.then(Vec::new),
53                ..ExecaResult::default()
54            },
55            options,
56            started: Instant::now(),
57            stopped: false,
58        }
59    }
60
61    pub fn pid(&self) -> Option<u32> {
62        self.stream.pid()
63    }
64
65    pub async fn wait_for_pid(&mut self) -> Option<u32> {
66        self.stream.wait_for_pid().await
67    }
68
69    pub fn kill(&mut self, signal: &str) {
70        if !self.stopped {
71            self.stopped = true;
72            self.result.killed = true;
73            self.stream.kill_with(signal);
74        }
75    }
76
77    /// Read a typed chunk and retain it for `wait()` unless buffering is off.
78    pub async fn next(&mut self) -> Option<OutputChunk> {
79        loop {
80            let deadline = if self.stopped { None } else { self.deadline };
81            tokio::select! {
82                chunk = self.stream.next() => {
83                    if let Some(chunk) = &chunk { self.capture(chunk); }
84                    return chunk;
85                }
86                _ = wait_deadline(deadline) => {
87                    self.result.timed_out = true;
88                    self.kill(&self.options.kill_signal.clone());
89                }
90                _ = wait_cancel(&mut self.options.cancel_signal), if !self.stopped => {
91                    self.result.is_canceled = true;
92                    self.kill(&self.options.kill_signal.clone());
93                }
94            }
95        }
96    }
97
98    fn capture(&mut self, chunk: &OutputChunk) {
99        let (data, destination) = match chunk {
100            OutputChunk::Stdout(data) => (data, &mut self.result.stdout),
101            OutputChunk::Stderr(data) => (data, &mut self.result.stderr),
102            OutputChunk::Exit(code) => {
103                self.result.exit_code = Some(*code);
104                return;
105            }
106        };
107        if !self.options.buffer {
108            return;
109        }
110        let remaining = self.options.max_buffer.saturating_sub(destination.len());
111        destination.extend_from_slice(&data[..data.len().min(remaining)]);
112        if let Some(all) = &mut self.result.all {
113            let remaining = self
114                .options
115                .max_buffer
116                .saturating_mul(2)
117                .saturating_sub(all.len());
118            all.extend_from_slice(&data[..data.len().min(remaining)]);
119        }
120        if data.len() > remaining {
121            self.result.is_max_buffer = true;
122            self.kill(&self.options.kill_signal.clone());
123        }
124    }
125
126    /// Drain remaining output and wait for the process and its output pumps.
127    pub async fn wait(mut self) -> ExecaOutcome {
128        while self.next().await.is_some() {}
129        match (&mut self.task).await {
130            Ok(Ok(())) => {}
131            Ok(Err(error)) => self.result.cause = Some(error.to_string()),
132            Err(error) => self.result.cause = Some(error.to_string()),
133        }
134        self.result.signal = self.stream.exit_signal();
135        if self.result.signal.is_some() {
136            self.result.exit_code = None;
137        }
138        self.result.duration = self.started.elapsed();
139        self.result.failed = self.result.exit_code != Some(0)
140            || self.result.killed
141            || self.result.timed_out
142            || self.result.is_canceled
143            || self.result.is_max_buffer
144            || self.result.cause.is_some();
145        if self.options.strip_final_newline {
146            strip_final_newline(&mut self.result.stdout);
147            strip_final_newline(&mut self.result.stderr);
148            if let Some(all) = &mut self.result.all {
149                strip_final_newline(all);
150            }
151        }
152        crate::trace::trace_lazy("Execa", || {
153            format!("Finished | failed={}", self.result.failed)
154        });
155        if self.options.reject && self.result.failed {
156            Err(ExecaError {
157                result: Box::new(self.result),
158            })
159        } else {
160            Ok(self.result)
161        }
162    }
163}
164
165async fn wait_deadline(deadline: Option<TokioInstant>) {
166    match deadline {
167        Some(deadline) => sleep_until(deadline).await,
168        None => pending().await,
169    }
170}
171
172async fn wait_cancel(signal: &mut Option<CancelSignal>) {
173    if let Some(signal) = signal {
174        if signal.0.wait_for(|canceled| *canceled).await.is_ok() {
175            return;
176        }
177    }
178    pending().await
179}