Skip to main content

command_stream/
stream.rs

1//! Streaming and async iteration support
2//!
3//! This module provides async streaming capabilities similar to JavaScript's
4//! async iterators and stream handling in `$.stream-utils.mjs`.
5//!
6//! It mirrors the JavaScript implementation's behavior for issue #155:
7//!
8//!   1. The stream yields an explicit `OutputChunk::Exit(code)` when the
9//!      process exits, so consumers can observe the exit code from inside the
10//!      loop.
11//!   2. The stream does not hang forever when the process has exited but a
12//!      grandchild keeps the stdio pipes open (the readers are drained with a
13//!      grace period and then aborted).
14//!   3. The process can be stopped from inside the loop via
15//!      [`OutputStream::kill`] / [`OutputStream::kill_with`], and abandoning the
16//!      stream (e.g. `break`) also stops the process.
17//!   4. The stop signal is configurable via
18//!      [`StreamingRunner::kill_signal`] (default `SIGTERM`), just like the
19//!      JavaScript `killSignal` option.
20//!
21//! ## Usage
22//!
23//! ```rust,no_run
24//! use command_stream::{StreamingRunner, OutputChunk};
25//!
26//! #[tokio::main]
27//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
28//!     let runner = StreamingRunner::new("yes hello");
29//!
30//!     // Stream output as it arrives
31//!     let mut stream = runner.stream();
32//!     let mut count = 0;
33//!     while let Some(chunk) = stream.next().await {
34//!         match chunk {
35//!             OutputChunk::Stdout(data) => {
36//!                 print!("{}", String::from_utf8_lossy(&data));
37//!                 count += 1;
38//!                 if count >= 5 {
39//!                     // Stop the process from inside the loop.
40//!                     stream.kill();
41//!                 }
42//!             }
43//!             OutputChunk::Stderr(data) => {
44//!                 eprint!("{}", String::from_utf8_lossy(&data));
45//!             }
46//!             OutputChunk::Exit(code) => {
47//!                 println!("Process exited with code: {}", code);
48//!                 break;
49//!             }
50//!         }
51//!     }
52//!
53//!     Ok(())
54//! }
55//! ```
56
57use std::collections::HashMap;
58use std::ffi::OsString;
59use std::path::PathBuf;
60use std::process::Stdio;
61use std::time::Duration;
62use tokio::io::BufReader;
63use tokio::process::Command;
64use tokio::sync::{mpsc, watch};
65use tokio::task::JoinHandle;
66
67use crate::signal::{
68    send_signal_to_process, signal_exit_code, Delivery, DEFAULT_KILL_GRACE_MS, DEFAULT_KILL_SIGNAL,
69};
70use crate::trace::trace_lazy;
71use crate::{CommandResult, Result};
72
73/// Default grace period (in milliseconds) to keep draining the stdio pipes
74/// after the process has exited before aborting any lingering readers. Mirrors
75/// the JavaScript `exitPumpGrace` default.
76const DEFAULT_EXIT_PUMP_GRACE_MS: u64 = 100;
77
78/// A chunk of output from a streaming process
79#[derive(Debug, Clone)]
80pub enum OutputChunk {
81    /// Stdout data
82    Stdout(Vec<u8>),
83    /// Stderr data
84    Stderr(Vec<u8>),
85    /// Process exit code
86    Exit(i32),
87}
88
89/// A streaming process runner that allows async iteration over output
90pub struct StreamingRunner {
91    command: StreamingCommand,
92    cwd: Option<PathBuf>,
93    env: Option<HashMap<String, String>>,
94    env_clear: bool,
95    prefer_local: crate::PreferLocal,
96    stdin_content: Option<Vec<u8>>,
97    kill_signal: String,
98    kill_grace_ms: u64,
99    exit_pump_grace_ms: u64,
100}
101
102#[derive(Clone)]
103enum StreamingCommand {
104    Shell(String),
105    Argv {
106        program: OsString,
107        args: Vec<OsString>,
108    },
109}
110
111impl StreamingRunner {
112    /// Create a streaming runner for a command string interpreted by the
113    /// platform shell.
114    pub fn new(command: impl Into<String>) -> Self {
115        Self::with_command(StreamingCommand::Shell(command.into()))
116    }
117
118    /// Create a streaming runner for an executable and exact argument vector.
119    ///
120    /// Unlike [`StreamingRunner::new`], this constructor bypasses the platform
121    /// shell. Argument boundaries are therefore preserved on every platform,
122    /// including Windows, without requiring shell-specific quoting.
123    pub fn from_argv<P, I, S>(program: P, args: I) -> Self
124    where
125        P: Into<OsString>,
126        I: IntoIterator<Item = S>,
127        S: Into<OsString>,
128    {
129        Self::with_command(StreamingCommand::Argv {
130            program: program.into(),
131            args: args.into_iter().map(Into::into).collect(),
132        })
133    }
134
135    fn with_command(command: StreamingCommand) -> Self {
136        StreamingRunner {
137            command,
138            cwd: None,
139            env: None,
140            env_clear: false,
141            prefer_local: crate::PreferLocal::Off,
142            stdin_content: None,
143            kill_signal: DEFAULT_KILL_SIGNAL.to_string(),
144            kill_grace_ms: DEFAULT_KILL_GRACE_MS,
145            exit_pump_grace_ms: DEFAULT_EXIT_PUMP_GRACE_MS,
146        }
147    }
148
149    /// Set the working directory
150    pub fn cwd(mut self, path: impl Into<PathBuf>) -> Self {
151        self.cwd = Some(path.into());
152        self
153    }
154
155    /// Set environment variables
156    pub fn env(mut self, env: HashMap<String, String>) -> Self {
157        self.env = Some(env);
158        self
159    }
160
161    /// Prefer executables from the command's working directory or explicit directories.
162    pub fn prefer_local(mut self, preference: crate::PreferLocal) -> Self {
163        self.prefer_local = preference;
164        self
165    }
166
167    /// Set stdin content
168    pub fn stdin(mut self, content: impl Into<String>) -> Self {
169        self.stdin_content = Some(content.into().into_bytes());
170        self
171    }
172
173    /// Set binary stdin without UTF-8 conversion.
174    pub fn stdin_bytes(mut self, content: impl Into<Vec<u8>>) -> Self {
175        self.stdin_content = Some(content.into());
176        self
177    }
178
179    /// Clear inherited environment variables before applying `env`.
180    pub fn clear_env(mut self, clear: bool) -> Self {
181        self.env_clear = clear;
182        self
183    }
184
185    /// Configure the signal used to stop the process when it is killed without
186    /// an explicit signal — i.e. [`OutputStream::kill`] or abandoning the
187    /// stream. Mirrors the JavaScript `killSignal` option (default `SIGTERM`).
188    ///
189    /// The reported exit code follows the conventional `128 + signal` mapping
190    /// (e.g. `SIGTERM` => 143, `SIGINT` => 130, `SIGKILL` => 137).
191    pub fn kill_signal(mut self, signal: impl Into<String>) -> Self {
192        self.kill_signal = signal.into();
193        self
194    }
195
196    /// Configure how long (in milliseconds) the child is given to handle the
197    /// kill signal before `SIGKILL` is sent. Mirrors the JavaScript `killGrace`
198    /// option (default 100ms).
199    ///
200    /// This is the window in which a child running its own `SIGTERM` handler
201    /// can shut down on its own terms. Set it to `0` to escalate immediately.
202    pub fn kill_grace_ms(mut self, ms: u64) -> Self {
203        self.kill_grace_ms = ms;
204        self
205    }
206
207    /// Configure the grace period (in milliseconds) to keep draining the stdio
208    /// pipes after the process exits before aborting lingering readers. Mirrors
209    /// the JavaScript `exitPumpGrace` option (default 100ms).
210    pub fn exit_pump_grace_ms(mut self, ms: u64) -> Self {
211        self.exit_pump_grace_ms = ms;
212        self
213    }
214
215    pub(crate) fn spawn(mut self) -> (OutputStream, JoinHandle<Result<()>>) {
216        let (tx, rx) = mpsc::channel(1024);
217        // Unbounded so a synchronous Drop can request a kill without awaiting.
218        let (kill_tx, kill_rx) = mpsc::unbounded_channel::<String>();
219        // The child is spawned inside the task below, so its id is not known
220        // when this returns. The task publishes it here as soon as the spawn
221        // succeeds; `OutputStream::pid` reads the latest value (issue #18).
222        let (pid_tx, pid_rx) = watch::channel(None);
223        let (exit_signal_tx, exit_signal_rx) = watch::channel(None);
224
225        // Spawn the process handling task
226        let command = self.command.clone();
227        let cwd = self.cwd.take();
228        let mut env = self.env.take();
229        let local_cwd = cwd
230            .clone()
231            .or_else(|| std::env::current_dir().ok())
232            .unwrap_or_else(|| PathBuf::from("."));
233        if let Some((key, path)) =
234            crate::local_bin::preferred_path(env.as_ref(), &local_cwd, &self.prefer_local)
235        {
236            env.get_or_insert_with(|| std::env::vars().collect())
237                .insert(key, path);
238        }
239        let stdin_content = self.stdin_content.take();
240        let env_clear = self.env_clear;
241        let grace = GraceWindows {
242            exit_pump_ms: self.exit_pump_grace_ms,
243            kill_ms: self.kill_grace_ms,
244        };
245        let kill_signal = self.kill_signal.clone();
246
247        let task = tokio::spawn(async move {
248            let channels = StreamChannels {
249                output_tx: tx,
250                kill_rx,
251                pid_tx,
252                exit_signal_tx,
253            };
254            let result =
255                run_streaming_process(command, cwd, env, stdin_content, env_clear, grace, channels)
256                    .await;
257            if let Err(error) = &result {
258                trace_lazy("StreamingRunner", || format!("Error: {error}"));
259            }
260            result
261        });
262
263        (
264            OutputStream {
265                rx,
266                kill_tx,
267                kill_signal,
268                killed: false,
269                pid_rx,
270                exit_signal_rx,
271            },
272            task,
273        )
274    }
275
276    /// Start the process and return a stream of output chunks
277    pub fn stream(self) -> OutputStream {
278        self.spawn().0
279    }
280
281    /// Run to completion and collect all output
282    pub async fn collect(self) -> Result<CommandResult> {
283        let stdin_content = self.stdin_content.clone();
284        let mut stdout = Vec::new();
285        let mut stderr = Vec::new();
286        let mut exit_code = 0;
287
288        let (mut stream, task) = self.spawn();
289        while let Some(chunk) = stream.rx.recv().await {
290            match chunk {
291                OutputChunk::Stdout(data) => stdout.extend(data),
292                OutputChunk::Stderr(data) => stderr.extend(data),
293                OutputChunk::Exit(code) => exit_code = code,
294            }
295        }
296
297        task.await.map_err(|error| {
298            std::io::Error::other(format!("streaming process task failed: {error}"))
299        })??;
300
301        let mut result = CommandResult::new(
302            String::from_utf8_lossy(&stdout).to_string(),
303            String::from_utf8_lossy(&stderr).to_string(),
304            exit_code,
305        );
306        if let Some(content) = stdin_content {
307            result.stdin = crate::result_streams::CapturedInput::new(content);
308        }
309        Ok(result)
310    }
311
312    /// Run an exact-argument command to completion from synchronous code.
313    ///
314    /// Build the command with [`Self::from_argv`] and configure it with the
315    /// same `cwd`, `env`, and `stdin` methods used by [`Self::collect`]. Call
316    /// this outside a Tokio runtime; async callers should use `collect().await`.
317    pub fn collect_blocking(self) -> Result<CommandResult> {
318        if tokio::runtime::Handle::try_current().is_ok() {
319            return Err(std::io::Error::other(
320                "collect_blocking cannot run inside a Tokio runtime; use collect().await",
321            )
322            .into());
323        }
324        let runtime = tokio::runtime::Runtime::new()?;
325        runtime.block_on(self.collect())
326    }
327}
328
329/// Stream of output chunks from a process
330pub struct OutputStream {
331    rx: mpsc::Receiver<OutputChunk>,
332    kill_tx: mpsc::UnboundedSender<String>,
333    kill_signal: String,
334    killed: bool,
335    pid_rx: watch::Receiver<Option<u32>>,
336    exit_signal_rx: watch::Receiver<Option<String>>,
337}
338
339impl OutputStream {
340    /// Receive the next chunk
341    pub async fn next(&mut self) -> Option<OutputChunk> {
342        self.rx.recv().await
343    }
344
345    /// Signal reported by the native child, available once it exits.
346    /// Ordinary exit codes above 128 are not mistaken for signals.
347    pub fn exit_signal(&self) -> Option<String> {
348        self.exit_signal_rx.borrow().clone()
349    }
350
351    /// Process id of the streamed command, as currently known.
352    ///
353    /// The child is spawned by a background task, so this is `None` for the
354    /// short window between [`StreamingRunner::stream`] returning and the spawn
355    /// completing, and stays `None` if the spawn failed. From the first
356    /// delivered chunk onwards it is set, and it remains readable after the
357    /// process has exited. Use [`wait_for_pid`](Self::wait_for_pid) to avoid
358    /// the startup window.
359    pub fn pid(&self) -> Option<u32> {
360        *self.pid_rx.borrow()
361    }
362
363    /// Process id of the streamed command, waiting for the spawn to complete.
364    ///
365    /// Resolves as soon as the child exists, and returns `None` if the process
366    /// could never be spawned. This is the streaming counterpart of awaiting a
367    /// stream before reading `runner.pid` in JavaScript.
368    pub async fn wait_for_pid(&mut self) -> Option<u32> {
369        // `wait_for` checks the current value first, so an already-published id
370        // returns without waiting. An error means the sending task is gone,
371        // which only happens when the spawn failed.
372        match self.pid_rx.wait_for(|pid| pid.is_some()).await {
373            Ok(pid) => *pid,
374            Err(_) => None,
375        }
376    }
377
378    /// Stop the process using the configured kill signal (default `SIGTERM`).
379    ///
380    /// This can be called from inside the consumption loop to stop a
381    /// long-running or endless process; a terminating `OutputChunk::Exit` is
382    /// still delivered afterwards.
383    pub fn kill(&mut self) {
384        let signal = self.kill_signal.clone();
385        self.kill_with(&signal);
386    }
387
388    /// Stop the process using an explicit signal, overriding the configured
389    /// kill signal for this call.
390    pub fn kill_with(&mut self, signal: &str) {
391        if self.killed {
392            return;
393        }
394        self.killed = true;
395        trace_lazy("OutputStream", || format!("kill | signal={}", signal));
396        // Best effort: the task may have already finished, in which case the
397        // receiver is gone and the send fails harmlessly.
398        let _ = self.kill_tx.send(signal.to_string());
399    }
400
401    /// Collect all remaining output into vectors
402    pub async fn collect(mut self) -> (Vec<u8>, Vec<u8>, i32) {
403        let mut stdout = Vec::new();
404        let mut stderr = Vec::new();
405        let mut exit_code = 0;
406
407        while let Some(chunk) = self.rx.recv().await {
408            match chunk {
409                OutputChunk::Stdout(data) => stdout.extend(data),
410                OutputChunk::Stderr(data) => stderr.extend(data),
411                OutputChunk::Exit(code) => exit_code = code,
412            }
413        }
414
415        (stdout, stderr, exit_code)
416    }
417
418    /// Collect stdout only, discarding stderr
419    pub async fn collect_stdout(mut self) -> Vec<u8> {
420        let mut stdout = Vec::new();
421
422        while let Some(chunk) = self.rx.recv().await {
423            if let OutputChunk::Stdout(data) = chunk {
424                stdout.extend(data);
425            }
426        }
427
428        stdout
429    }
430}
431
432impl Drop for OutputStream {
433    fn drop(&mut self) {
434        // Abandoning the stream (e.g. `break`-ing out of the loop) must stop the
435        // process, matching the JavaScript iterator's `finally` cleanup. If the
436        // process already finished this is a harmless no-op.
437        if !self.killed {
438            let _ = self.kill_tx.send(self.kill_signal.clone());
439        }
440    }
441}
442
443/// The channels `run_streaming_process` communicates over: output chunks out,
444/// kill requests in, and the child's id published once the spawn succeeds.
445struct StreamChannels {
446    /// Carries the output chunks, and finally the `Exit` chunk, to the consumer.
447    output_tx: mpsc::Sender<OutputChunk>,
448    /// Carries kill requests, by signal name, in from the consumer.
449    kill_rx: mpsc::UnboundedReceiver<String>,
450    /// Publishes the child's id, which is only known inside the spawning task.
451    pid_tx: watch::Sender<Option<u32>>,
452    exit_signal_tx: watch::Sender<Option<String>>,
453}
454
455/// How long the runner waits, in milliseconds, at the two points where it gives
456/// something a chance to finish on its own before forcing the issue.
457#[derive(Debug, Clone, Copy)]
458struct GraceWindows {
459    /// Time allowed for the readers to drain buffered output after the child
460    /// exits, before the `Exit` chunk is emitted.
461    exit_pump_ms: u64,
462    /// Time allowed for the child to handle the delivered signal, before the
463    /// escalation to `SIGKILL`.
464    kill_ms: u64,
465}
466
467/// Run a streaming process and send output to the channel
468async fn run_streaming_process(
469    command: StreamingCommand,
470    cwd: Option<PathBuf>,
471    env: Option<HashMap<String, String>>,
472    stdin_content: Option<Vec<u8>>,
473    env_clear: bool,
474    grace: GraceWindows,
475    channels: StreamChannels,
476) -> Result<()> {
477    let StreamChannels {
478        output_tx: tx,
479        mut kill_rx,
480        pid_tx,
481        exit_signal_tx,
482    } = channels;
483    trace_lazy("StreamingRunner", || match &command {
484        StreamingCommand::Shell(command) => format!("Starting: {command}"),
485        StreamingCommand::Argv { program, args } => {
486            format!("Starting argv command: {program:?} {args:?}")
487        }
488    });
489
490    let mut cmd = match command {
491        StreamingCommand::Shell(command) => crate::utils::shell_command(&command, env.as_ref()),
492        StreamingCommand::Argv { program, args } => {
493            let mut cmd = Command::new(program);
494            cmd.args(args);
495            cmd
496        }
497    };
498
499    cmd.kill_on_drop(true);
500
501    // Configure stdio
502    if stdin_content.is_some() {
503        cmd.stdin(Stdio::piped());
504    } else {
505        cmd.stdin(Stdio::null());
506    }
507    cmd.stdout(Stdio::piped());
508    cmd.stderr(Stdio::piped());
509
510    // Run the child in its own process group so we can signal the whole group
511    // (parent + grandchildren), matching the JavaScript implementation.
512    #[cfg(unix)]
513    cmd.process_group(0);
514
515    // Set working directory
516    if let Some(ref cwd) = cwd {
517        cmd.current_dir(cwd);
518    }
519
520    // Set environment
521    if env_clear {
522        cmd.env_clear();
523    }
524    if let Some(ref env_vars) = env {
525        for (key, value) in env_vars {
526            cmd.env(key, value);
527        }
528    }
529
530    let mut child = cmd.spawn()?;
531    // Publish the id before any awaiting, so a consumer asking for it as soon
532    // as the first chunk arrives already sees it.
533    let _ = pid_tx.send(child.id());
534
535    // Write stdin concurrently with output readers and process cancellation.
536    // Filling stdin before reading stdout deadlocks a full-duplex child.
537    let stdin_handle = stdin_content.and_then(|content| {
538        child.stdin.take().map(|mut stdin| {
539            tokio::spawn(async move {
540                use tokio::io::AsyncWriteExt;
541                let _ = stdin.write_all(&content).await;
542                let _ = stdin.shutdown().await;
543            })
544        })
545    });
546
547    // Spawn stdout reader
548    let stdout = child.stdout.take();
549    let tx_stdout = tx.clone();
550    let stdout_handle = stdout.map(|stdout| {
551        tokio::spawn(async move {
552            let mut reader = BufReader::new(stdout);
553            let mut buf = vec![0u8; 8192];
554            loop {
555                use tokio::io::AsyncReadExt;
556                match reader.read(&mut buf).await {
557                    Ok(0) => break,
558                    Ok(n) => {
559                        if tx_stdout
560                            .send(OutputChunk::Stdout(buf[..n].to_vec()))
561                            .await
562                            .is_err()
563                        {
564                            break;
565                        }
566                    }
567                    Err(_) => break,
568                }
569            }
570        })
571    });
572
573    // Spawn stderr reader
574    let stderr = child.stderr.take();
575    let tx_stderr = tx.clone();
576    let stderr_handle = stderr.map(|stderr| {
577        tokio::spawn(async move {
578            let mut reader = BufReader::new(stderr);
579            let mut buf = vec![0u8; 8192];
580            loop {
581                use tokio::io::AsyncReadExt;
582                match reader.read(&mut buf).await {
583                    Ok(0) => break,
584                    Ok(n) => {
585                        if tx_stderr
586                            .send(OutputChunk::Stderr(buf[..n].to_vec()))
587                            .await
588                            .is_err()
589                        {
590                            break;
591                        }
592                    }
593                    Err(_) => break,
594                }
595            }
596        })
597    });
598
599    // Wait for the process to exit OR for a kill request — crucially we do NOT
600    // wait for the readers first. If a grandchild keeps the pipe open the
601    // readers would never finish, so waiting on them before the exit (as the
602    // old implementation did) would hang forever (issue #155).
603    let pid = child.id();
604    let code;
605    tokio::select! {
606        status = child.wait() => {
607            let status = status?;
608            #[cfg(unix)]
609            {
610                use std::os::unix::process::ExitStatusExt;
611                let signal = status.signal().and_then(|number| {
612                    nix::sys::signal::Signal::try_from(number).ok().map(|signal| signal.to_string())
613                });
614                let _ = exit_signal_tx.send(signal);
615            }
616            code = status_to_code(status);
617        }
618        maybe_signal = kill_rx.recv() => {
619            // A kill was requested (explicit kill()/kill_with() or the stream
620            // being dropped). Stop the process group with the requested signal.
621            let signal = maybe_signal.unwrap_or_else(|| DEFAULT_KILL_SIGNAL.to_string());
622            trace_lazy("StreamingRunner", || format!("Kill requested | signal={}", signal));
623            // Give the child its grace period to run its own handler and exit
624            // on its own terms, then escalate to a forceful kill so a process
625            // that ignores the signal still terminates.
626            //
627            // A zero grace period means the child is given no opportunity to
628            // handle the signal, so the requested signal is not delivered at
629            // all. Anything done between it and the forceful kill - a syscall,
630            // or awaiting a zero-length timeout, which yields to the runtime -
631            // is a window the child can be scheduled in, which made "no grace"
632            // a race the child occasionally won rather than a guarantee.
633            if cfg!(unix) && grace.kill_ms > 0 && signal != "SIGKILL" {
634                if let Some(pid) = pid {
635                    // The child is always spawned with `process_group(0)`
636                    // above, so it leads the group named by its own pid.
637                    send_signal_to_process(pid, &signal, Delivery::ProcessAndGroup);
638                }
639                // Keep the leader unreaped until escalation: its descendants
640                // may ignore the first signal even when the leader exits.
641                // Reaping it here previously suppressed their SIGKILL.
642                tokio::time::sleep(Duration::from_millis(grace.kill_ms)).await;
643            }
644            if let Some(pid) = pid {
645                // Windows uses taskkill here before terminating the parent.
646                send_signal_to_process(pid, "SIGKILL", Delivery::ProcessAndGroup);
647            }
648            let _ = child.start_kill();
649            let _ = child.wait().await;
650            // Report the conventional 128 + signal code for the requested
651            // signal, matching the JavaScript implementation.
652            let _ = exit_signal_tx.send(Some(signal.clone()));
653            code = signal_exit_code(&signal);
654        }
655    }
656
657    if let Some(handle) = stdin_handle {
658        handle.abort();
659        let _ = handle.await;
660    }
661
662    // The process has exited. Give the readers a short grace period to flush any
663    // buffered output, then abort any that are still blocked on an inherited
664    // open pipe so we don't hang.
665    let stdout_abort = stdout_handle.as_ref().map(|h| h.abort_handle());
666    let stderr_abort = stderr_handle.as_ref().map(|h| h.abort_handle());
667    let drain = async {
668        if let Some(handle) = stdout_handle {
669            let _ = handle.await;
670        }
671        if let Some(handle) = stderr_handle {
672            let _ = handle.await;
673        }
674    };
675    if tokio::time::timeout(Duration::from_millis(grace.exit_pump_ms), drain)
676        .await
677        .is_err()
678    {
679        // A reader is still blocked on an inherited open pipe — abort it so the
680        // exit chunk is delivered without waiting for the grandchild.
681        if let Some(abort) = stdout_abort {
682            abort.abort();
683        }
684        if let Some(abort) = stderr_abort {
685            abort.abort();
686        }
687    }
688
689    // Send exit code (always — even if a reader was aborted).
690    let _ = tx.send(OutputChunk::Exit(code)).await;
691
692    trace_lazy("StreamingRunner", || format!("Exited with code: {}", code));
693
694    Ok(())
695}
696
697/// Convert an exit status into a numeric exit code, using the conventional
698/// `128 + signal` mapping when the process was terminated by a signal.
699fn status_to_code(status: std::process::ExitStatus) -> i32 {
700    if let Some(code) = status.code() {
701        return code;
702    }
703    #[cfg(unix)]
704    {
705        use std::os::unix::process::ExitStatusExt;
706        if let Some(sig) = status.signal() {
707            return 128 + sig;
708        }
709    }
710    -1
711}
712
713/// Async iterator trait for output streams
714#[async_trait::async_trait]
715pub trait AsyncIterator {
716    type Item;
717
718    /// Get the next item from the iterator
719    async fn next(&mut self) -> Option<Self::Item>;
720}
721
722#[async_trait::async_trait]
723impl AsyncIterator for OutputStream {
724    type Item = OutputChunk;
725
726    async fn next(&mut self) -> Option<Self::Item> {
727        self.rx.recv().await
728    }
729}
730
731/// Extension trait to convert ProcessRunner into a stream
732pub trait IntoStream {
733    /// Convert into an output stream
734    fn into_stream(self) -> OutputStream;
735}
736
737impl IntoStream for crate::ProcessRunner {
738    fn into_stream(self) -> OutputStream {
739        let streaming = StreamingRunner::new(self.command().to_string());
740        streaming.stream()
741    }
742}