command_stream/execa/
process.rs1use 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
9pub 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 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 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}