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#[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 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#[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: Option<u32>,
231 span: Span,
232}
233
234#[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 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 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 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 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}