Skip to main content

command_stream/
lib.rs

1//! # command-stream
2//!
3//! Modern shell command execution library with streaming, async iteration, and event support.
4//!
5//! This library provides a Rust equivalent to the JavaScript command-stream library,
6//! offering powerful shell command execution with streaming capabilities.
7//!
8//! ## Features
9//!
10//! - Async command execution with tokio
11//! - Streaming output via async iterators
12//! - Event-based output handling (on, once, emit)
13//! - Virtual commands for common operations (cat, ls, mkdir, etc.)
14//! - Shell operator support (&&, ||, ;, |)
15//! - Pipeline support with `.pipe()` method and `Pipeline` builder
16//! - Global state management for shell settings
17//! - `cmd!` macro for ergonomic command creation (similar to JS `$` tagged template literals)
18//! - Cross-platform support
19//!
20//! ## Module Organization
21//!
22//! The codebase follows a modular architecture similar to the JavaScript implementation:
23//!
24//! - `ansi` - ANSI escape code handling utilities
25//! - `commands` - Virtual command implementations
26//! - `events` - Event emitter for stream events
27//! - `macros` - The `cmd!` macro for ergonomic command creation
28//! - `pipeline` - Pipeline execution support
29//! - `quote` - Shell quoting utilities
30//! - `shell_parser` - Shell command parsing
31//! - `state` - Global state management
32//! - `stream` - Async streaming and iteration support
33//! - `trace` - Logging and tracing utilities
34//! - `utils` - Command results and virtual command helpers
35//!
36//! ## Quick Start
37//!
38//! ```rust,no_run
39//! use command_stream::{run, cmd};
40//!
41//! #[tokio::main]
42//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
43//!     // Execute a simple command
44//!     let result = run("echo hello world").await?;
45//!     println!("{}", result.stdout);
46//!
47//!     // Using the cmd! macro (similar to JS $ tagged template)
48//!     let name = "world";
49//!     let result = cmd!("echo hello {}", name).await?;
50//!     println!("{}", result.stdout);
51//!
52//!     // Using pipelines
53//!     use command_stream::Pipeline;
54//!     let result = Pipeline::new()
55//!         .add("echo hello world")
56//!         .add("grep world")
57//!         .run()
58//!         .await?;
59//!
60//!     Ok(())
61//! }
62//! ```
63
64// Modular utility modules (following JavaScript modular pattern)
65pub mod ansi;
66pub mod bun_shell;
67pub mod events;
68pub mod execa;
69pub mod local_bin;
70#[doc(hidden)]
71pub mod macros;
72pub mod pipeline;
73pub mod quote;
74pub mod result_streams;
75mod run_options;
76pub mod signal;
77pub mod state;
78pub mod stream;
79pub mod terminal;
80pub mod trace;
81
82// Core modules
83pub mod commands;
84pub mod shell_parser;
85pub mod utils;
86
87// zx-compatible API (google/zx 8.x): `$`-style Shell, ProcessOutput, goods.
88pub mod zx;
89
90use std::path::PathBuf;
91use std::process::Stdio;
92use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt};
93use tokio::process::Child;
94use tokio::sync::mpsc;
95
96pub use commands::{CommandContext, StreamChunk};
97pub use execa::{execa, execa_compat, execa_node, execa_sync};
98pub use run_options::{Error, Result, RunOptions, StdinOption};
99pub use shell_parser::{needs_real_shell, parse_shell_command, ParsedCommand};
100pub use utils::{CommandResult, VirtualUtils};
101
102// Re-export modular utilities at crate root for convenient access
103pub use ansi::{AnsiConfig, AnsiUtils};
104pub use events::{EventData, EventType, StreamEmitter};
105pub use local_bin::PreferLocal;
106pub use pipeline::{Pipeline, PipelineBuilder, PipelineExt};
107pub use quote::{
108    escape_for_double_quotes, escape_for_single_quotes, has_shell_escapes,
109    is_pre_quoted_passthrough_enabled, is_quote_context_enabled, quote, quote_for_context,
110    scan_quote_context, QuoteContext,
111};
112pub use signal::{signal_exit_code, signal_number, DEFAULT_KILL_GRACE_MS, DEFAULT_KILL_SIGNAL};
113pub use state::{
114    get_shell_settings, global_state, reset_global_state, set_shell_option, unset_shell_option,
115    GlobalState, ShellSettings,
116};
117pub use stream::{AsyncIterator, IntoStream, OutputChunk, OutputStream, StreamingRunner};
118pub use trace::trace;
119
120#[derive(Clone, Copy)]
121enum ChildOutput {
122    Stdout,
123    Stderr,
124}
125
126/// Read child output as byte chunks so capture does not invent a trailing newline.
127///
128/// stdout and stderr use separate futures in `ProcessRunner::run`, preventing
129/// either pipe from filling while the other is being drained. Mirroring keeps
130/// the original bytes too, including output that does not end in a newline.
131async fn collect_child_output<R>(
132    reader: Option<R>,
133    mirror: bool,
134    target: ChildOutput,
135) -> std::io::Result<Vec<u8>>
136where
137    R: AsyncRead + Unpin,
138{
139    let Some(mut reader) = reader else {
140        return Ok(Vec::new());
141    };
142    let mut collected = Vec::new();
143    let mut buffer = [0_u8; 8192];
144
145    loop {
146        let count = reader.read(&mut buffer).await?;
147        if count == 0 {
148            break;
149        }
150
151        let chunk = &buffer[..count];
152        collected.extend_from_slice(chunk);
153        if mirror {
154            match target {
155                ChildOutput::Stdout => {
156                    let mut output = std::io::stdout().lock();
157                    let _ = std::io::Write::write_all(&mut output, chunk);
158                    let _ = std::io::Write::flush(&mut output);
159                }
160                ChildOutput::Stderr => {
161                    let mut output = std::io::stderr().lock();
162                    let _ = std::io::Write::write_all(&mut output, chunk);
163                    let _ = std::io::Write::flush(&mut output);
164                }
165            }
166        }
167    }
168
169    Ok(collected)
170}
171
172fn fallback_cwd() -> PathBuf {
173    std::env::var_os("HOME")
174        .or_else(|| std::env::var_os("USERPROFILE"))
175        .map(PathBuf::from)
176        .filter(|path| path.is_dir())
177        .unwrap_or_else(std::env::temp_dir)
178}
179
180/// Resolve a working directory that is safe to spawn a child process in.
181///
182/// When no explicit cwd is requested the child normally inherits the parent's
183/// working directory. But if that directory has been deleted or become
184/// inaccessible (the "getcwd() failed" scenario from issue #44), inheriting it
185/// makes the OS-level spawn fail. In that case fall back to a directory that is
186/// known to exist so the command still runs.
187///
188/// Normal behavior is preserved: when an explicit cwd is given, or when the
189/// inherited working directory is valid, this returns the requested value
190/// (`None` meaning "inherit").
191fn resolve_spawn_cwd(cwd: Option<&PathBuf>) -> Option<PathBuf> {
192    // An explicit directory is always honored as-is.
193    if let Some(c) = cwd {
194        return Some(c.clone());
195    }
196
197    // No explicit cwd: we would inherit the parent's working directory. Make
198    // sure that directory is actually usable before relying on inheritance.
199    match std::env::current_dir() {
200        Ok(_) => None,
201        Err(e) => {
202            let fallback = fallback_cwd();
203            trace(
204                "ProcessRunner",
205                &format!(
206                    "current_dir() failed ({}); spawning in fallback directory {}",
207                    e,
208                    fallback.display()
209                ),
210            );
211            Some(fallback)
212        }
213    }
214}
215
216/// A running or completed process
217pub struct ProcessRunner {
218    command: String,
219    options: RunOptions,
220    child: Option<Child>,
221    stdin_bytes: Vec<u8>,
222    /// Process id of the spawned child, recorded at spawn time. `run()` takes
223    /// the child in order to await it, so reading the id from it only works
224    /// between `start()` and `run()`; this copy is what makes `pid()` answer
225    /// after the command has finished too (issue #18).
226    pid: Option<u32>,
227    result: Option<CommandResult>,
228    started: bool,
229    finished: bool,
230    cancelled: bool,
231    /// Whether the child was spawned into a process group of its own, and so
232    /// can be signalled as a group. Recorded at spawn time because it cannot be
233    /// discovered later: by the time the group is signalled the leader is
234    /// usually a zombie, which macOS refuses to answer `getpgid` for.
235    #[cfg(unix)]
236    own_process_group: bool,
237    output_tx: Option<mpsc::Sender<StreamChunk>>,
238    // Held, never read: dropping the receiver would close the channel, and
239    // streaming virtual commands treat a closed channel as "stop now" (see
240    // `commands::yes`, which loops until `output_tx.send` fails). Keeping it
241    // alive is what gives those commands their run-until-cancelled behaviour.
242    #[allow(dead_code)]
243    output_rx: Option<mpsc::Receiver<StreamChunk>>,
244}
245
246/// Borrowed access to the operating-system child owned by a [`ProcessRunner`].
247///
248/// The wrapper keeps process termination on the runner's signal-aware path:
249/// [`kill`](Self::kill) and [`kill_with`](Self::kill_with) signal the child and
250/// its process group, honor the configured grace period, and then escalate if
251/// necessary. Use [`native`](Self::native) or [`native_mut`](Self::native_mut)
252/// when direct access to Tokio's child object is required.
253pub struct ProcessChild<'a> {
254    runner: &'a mut ProcessRunner,
255}
256
257impl ProcessChild<'_> {
258    /// Process id of the active child.
259    pub fn pid(&self) -> Option<u32> {
260        self.native().id()
261    }
262
263    /// Borrow Tokio's native child process object.
264    pub fn native(&self) -> &Child {
265        self.runner
266            .child
267            .as_ref()
268            .expect("ProcessChild exists only while its native child is present")
269    }
270
271    /// Mutably borrow Tokio's native child process object.
272    pub fn native_mut(&mut self) -> &mut Child {
273        self.runner
274            .child
275            .as_mut()
276            .expect("ProcessChild exists only while its native child is present")
277    }
278
279    /// Stop the child using the runner's configured signal and grace period.
280    pub fn kill(&mut self) -> Result<()> {
281        self.runner.kill()
282    }
283
284    /// Stop the child using an explicit signal and the configured grace period.
285    pub fn kill_with(&mut self, signal: &str) -> Result<()> {
286        self.runner.kill_with(signal)
287    }
288}
289
290impl ProcessRunner {
291    /// Create a new process runner
292    pub fn new(command: impl Into<String>, options: RunOptions) -> Self {
293        let (tx, rx) = mpsc::channel(1024);
294        ProcessRunner {
295            command: command.into(),
296            options,
297            child: None,
298            stdin_bytes: Vec::new(),
299            pid: None,
300            result: None,
301            started: false,
302            finished: false,
303            cancelled: false,
304            #[cfg(unix)]
305            own_process_group: false,
306            output_tx: Some(tx),
307            output_rx: Some(rx),
308        }
309    }
310
311    /// Whether the child will read from the caller's terminal.
312    ///
313    /// Only an *inherited* stdin that is actually a tty counts: a pipe, a null
314    /// stdin, or inherited stdin that has been redirected to a file carries no
315    /// terminal, and neither does output-only inheritance. This is the one case
316    /// where the child must stay in the caller's process group.
317    #[cfg(unix)]
318    fn shares_the_terminal(&self) -> bool {
319        use std::io::IsTerminal;
320
321        self.options.interactive
322            || (matches!(self.options.stdin, StdinOption::Inherit)
323                && std::io::stdin().is_terminal())
324    }
325
326    /// Start the process
327    pub async fn start(&mut self) -> Result<()> {
328        if self.started {
329            return Ok(());
330        }
331        self.started = true;
332
333        utils::trace_lazy("ProcessRunner", || {
334            format!("Starting command: {}", self.command)
335        });
336
337        // Check if this is a virtual command. Backslash escapes are removed by a
338        // real shell but not by the whitespace splitting used for virtual
339        // command args, so such commands always go to the system shell (#49).
340        // The same applies to redirection and expansions: whitespace splitting
341        // would hand `>`, `out.txt` to the virtual command as two ordinary
342        // arguments, so `echo hello > out.txt` printed the redirection instead
343        // of writing the file, and `git push ... 2>&1` reported success while
344        // nothing was pushed (#46).
345        let first_word = if matches!(self.options.stdin, StdinOption::Pipe)
346            || has_shell_escapes(&self.command)
347            || needs_real_shell(&self.command)
348        {
349            ""
350        } else {
351            self.command.split_whitespace().next().unwrap_or("")
352        };
353        if let Some(mut result) = self.try_virtual_command(first_word).await {
354            if let StdinOption::Content(ref content) = self.options.stdin {
355                result.stdin =
356                    crate::result_streams::CapturedInput::new(content.as_bytes().to_vec());
357            }
358            self.result = Some(result);
359            self.finished = true;
360            return Ok(());
361        }
362        // Parse command for shell operators (for future use with virtual command pipelines)
363        let _parsed = if self.options.shell_operators && !needs_real_shell(&self.command) {
364            parse_shell_command(&self.command)
365        } else {
366            None
367        };
368
369        // Execute via real shell if needed
370        let mut cmd = utils::shell_command(&self.command, self.options.env.as_ref());
371
372        // Configure stdin
373        match &self.options.stdin {
374            StdinOption::Inherit => {
375                cmd.stdin(Stdio::inherit());
376            }
377            StdinOption::Pipe => {
378                cmd.stdin(Stdio::piped());
379            }
380            StdinOption::Content(_) => {
381                cmd.stdin(Stdio::piped());
382            }
383            StdinOption::Null => {
384                cmd.stdin(Stdio::null());
385            }
386        }
387
388        // Configure stdout/stderr
389        if self.options.capture || self.options.mirror {
390            cmd.stdout(Stdio::piped());
391            cmd.stderr(Stdio::piped());
392        } else {
393            cmd.stdout(Stdio::inherit());
394            cmd.stderr(Stdio::inherit());
395        }
396
397        // Set working directory. Fall back to a valid directory when the
398        // inherited working directory has been deleted (issue #44).
399        if let Some(cwd) = resolve_spawn_cwd(self.options.cwd.as_ref()) {
400            cmd.current_dir(cwd);
401        }
402
403        // Set environment
404        if let Some(ref env_vars) = self.options.env {
405            for (key, value) in env_vars {
406                cmd.env(key, value);
407            }
408        }
409        let local_cwd = self
410            .options
411            .cwd
412            .clone()
413            .or_else(|| std::env::current_dir().ok())
414            .unwrap_or_else(|| PathBuf::from("."));
415        local_bin::apply_prefer_local(
416            &mut cmd,
417            self.options.env.as_ref(),
418            &local_cwd,
419            &self.options.prefer_local,
420        );
421
422        // Run the child in its own process group so that killing it can signal
423        // the whole group (parent + grandchildren), matching `StreamingRunner`
424        // and the JavaScript implementation's `detached` spawn.
425        //
426        // A child that shares the terminal is deliberately left in the caller's
427        // group. The tty delivers CTRL+C to its foreground group only, so
428        // moving such a child out would both hide CTRL+C from it and stop it
429        // with SIGTTIN the moment it read from the terminal. JavaScript draws
430        // the same line, spawning interactive commands without `detached`.
431        #[cfg(unix)]
432        {
433            self.own_process_group = !self.shares_the_terminal();
434            if self.own_process_group {
435                cmd.process_group(0);
436            }
437        }
438
439        // Spawn the process
440        let child = cmd.spawn()?;
441        // Record the id while the child is still held. `run()` takes the child
442        // in order to await it, so this copy is what keeps `pid()` readable
443        // afterwards.
444        self.pid = child.id();
445        self.child = Some(child);
446
447        Ok(())
448    }
449
450    /// Borrow the active operating-system child.
451    ///
452    /// Call [`start`](Self::start) first. The result is `None` before startup,
453    /// for built-in commands (which run in-process), and after [`run`](Self::run)
454    /// consumes and reaps the child. Killing through the returned handle keeps
455    /// the runner's process-group and graceful-escalation behavior.
456    ///
457    /// ```no_run
458    /// use command_stream::{ProcessRunner, RunOptions};
459    ///
460    /// # #[tokio::main]
461    /// # async fn main() -> command_stream::Result<()> {
462    /// let mut runner = ProcessRunner::new("sleep 30", RunOptions::default());
463    /// runner.start().await?;
464    /// if let Some(mut child) = runner.child() {
465    ///     println!("child pid: {:?}", child.pid());
466    ///     child.kill_with("SIGTERM")?;
467    /// }
468    /// let _ = runner.run().await?;
469    /// # Ok(())
470    /// # }
471    /// ```
472    pub fn child(&mut self) -> Option<ProcessChild<'_>> {
473        self.child.as_ref()?;
474        Some(ProcessChild { runner: self })
475    }
476
477    /// Write bytes to the stdin pipe of a running command.
478    ///
479    /// Configure the runner with [`StdinOption::Pipe`], call [`start`](Self::start),
480    /// write as many chunks as needed, and finish with [`close_stdin`](Self::close_stdin).
481    pub async fn write_stdin(&mut self, data: impl AsRef<[u8]>) -> Result<()> {
482        self.start().await?;
483        let stdin = self
484            .child
485            .as_mut()
486            .and_then(|child| child.stdin.as_mut())
487            .ok_or_else(|| {
488                Error::Io(std::io::Error::new(
489                    std::io::ErrorKind::BrokenPipe,
490                    "command stdin is not available; use StdinOption::Pipe",
491                ))
492            })?;
493        stdin.write_all(data.as_ref()).await?;
494        self.stdin_bytes.extend_from_slice(data.as_ref());
495        Ok(())
496    }
497
498    /// Close a running command's stdin pipe so it can observe end-of-input.
499    pub async fn close_stdin(&mut self) -> Result<()> {
500        self.start().await?;
501        if let Some(mut stdin) = self.child.as_mut().and_then(|child| child.stdin.take()) {
502            stdin.shutdown().await?;
503        }
504        Ok(())
505    }
506
507    /// Run the process to completion
508    pub async fn run(&mut self) -> Result<CommandResult> {
509        self.start().await?;
510
511        if let Some(result) = &self.result {
512            return Ok(result.clone());
513        }
514
515        let mut child = self
516            .child
517            .take()
518            .ok_or_else(|| Error::Io(std::io::Error::other("Process not started")))?;
519
520        // Handle stdin content if provided
521        if let StdinOption::Content(ref content) = self.options.stdin {
522            if let Some(mut stdin) = child.stdin.take() {
523                let content = content.clone();
524                tokio::spawn(async move {
525                    let _ = stdin.write_all(content.as_bytes()).await;
526                    let _ = stdin.shutdown().await;
527                });
528            }
529        }
530
531        // Drain both pipes concurrently and preserve their newline framing. The
532        // previous line reader appended `\n` to every final line, changing
533        // output from commands such as `printf` that omit a newline (issue #37).
534        let stdout = child.stdout.take();
535        let stderr = child.stderr.take();
536        let collected = tokio::try_join!(
537            collect_child_output(stdout, self.options.mirror, ChildOutput::Stdout),
538            collect_child_output(stderr, self.options.mirror, ChildOutput::Stderr),
539        );
540        let (stdout, stderr) = match collected {
541            Ok(output) => output,
542            Err(error) => {
543                // `try_join!` drops the other pipe reader after an error. Stop
544                // and reap the child so it cannot remain blocked on that pipe.
545                let _ = child.start_kill();
546                let _ = child.wait().await;
547                return Err(error.into());
548            }
549        };
550
551        let status = child.wait().await?;
552        let code = status.code().unwrap_or(-1);
553
554        let mut result = CommandResult::new(
555            String::from_utf8_lossy(&stdout).into_owned(),
556            String::from_utf8_lossy(&stderr).into_owned(),
557            code,
558        );
559        if let StdinOption::Content(ref content) = self.options.stdin {
560            result.stdin = crate::result_streams::CapturedInput::new(content.as_bytes().to_vec());
561        } else if !self.stdin_bytes.is_empty() {
562            result.stdin = crate::result_streams::CapturedInput::new(self.stdin_bytes.clone());
563        }
564
565        self.result = Some(result.clone());
566        self.finished = true;
567
568        Ok(result)
569    }
570
571    /// Try to execute as a virtual command
572    async fn try_virtual_command(&self, cmd_name: &str) -> Option<CommandResult> {
573        if !commands::are_virtual_commands_enabled() {
574            return None;
575        }
576
577        // An empty command name means the caller already decided this command
578        // must go to a real shell (redirection, expansions, escapes). Bail out
579        // before tokenizing so we neither waste work nor parse shell syntax we
580        // deliberately delegate.
581        if cmd_name.is_empty() {
582            return None;
583        }
584
585        // Parse args from command string, respecting quotes and performing
586        // POSIX quote removal so `echo label:'help wanted'` reaches the built-in
587        // as the single argument `label:help wanted` (issue #48).
588        let words = shell_parser::split_command_words(&self.command);
589        let args: Vec<String> = words.into_iter().skip(1).collect();
590
591        let ctx = CommandContext {
592            args,
593            stdin: match &self.options.stdin {
594                StdinOption::Content(s) => Some(s.clone()),
595                _ => None,
596            },
597            cwd: self.options.cwd.clone(),
598            env: self.options.env.clone(),
599            output_tx: self.output_tx.clone(),
600            is_cancelled: None,
601        };
602
603        match cmd_name {
604            "echo" => Some(commands::echo(ctx).await),
605            "pwd" => Some(commands::pwd(ctx).await),
606            "cd" => Some(commands::cd::resolve_cd(ctx).await.0),
607            "true" => Some(commands::r#true(ctx).await),
608            "false" => Some(commands::r#false(ctx).await),
609            "sleep" => Some(commands::sleep(ctx).await),
610            "cat" => Some(commands::cat(ctx).await),
611            "ls" => Some(commands::ls(ctx).await),
612            "mkdir" => Some(commands::mkdir(ctx).await),
613            "rm" => Some(commands::rm(ctx).await),
614            "touch" => Some(commands::touch(ctx).await),
615            "cp" => Some(commands::cp(ctx).await),
616            "mv" => Some(commands::mv(ctx).await),
617            "basename" => Some(commands::basename(ctx).await),
618            "dirname" => Some(commands::dirname(ctx).await),
619            "env" => Some(commands::env(ctx).await),
620            "exit" => Some(commands::exit(ctx).await),
621            "which" => Some(commands::which(ctx).await),
622            "yes" => Some(commands::yes(ctx).await),
623            "seq" => Some(commands::seq(ctx).await),
624            "tee" => Some(commands::tee(ctx).await),
625            "test" => Some(commands::test(ctx).await),
626            _ => None,
627        }
628    }
629
630    /// Stop the process using the configured kill signal
631    /// ([`RunOptions::kill_signal`], default `SIGTERM`).
632    ///
633    /// Mirrors the JavaScript `kill()` with no argument.
634    pub fn kill(&mut self) -> Result<()> {
635        let signal = self.options.kill_signal.clone();
636        self.kill_with(&signal)
637    }
638
639    /// Stop the process using an explicit signal, overriding
640    /// [`RunOptions::kill_signal`] for this call.
641    ///
642    /// Mirrors the JavaScript `kill(signal)`. The signal is delivered to the
643    /// child and its process group, so grandchildren behind a `sh -c` wrapper
644    /// are stopped too - except for a child sharing the caller's terminal,
645    /// which stays in the caller's group so CTRL+C keeps reaching it. The child
646    /// then has [`RunOptions::kill_grace_ms`] to run its own handler before
647    /// `SIGKILL` follows, so a process that ignores the signal still
648    /// terminates.
649    ///
650    /// ```no_run
651    /// use command_stream::{ProcessRunner, RunOptions};
652    ///
653    /// # #[tokio::main]
654    /// # async fn main() -> command_stream::Result<()> {
655    /// let mut runner = ProcessRunner::new("sleep 30", RunOptions::default());
656    /// runner.start().await?;
657    /// runner.kill_with("SIGINT")?; // the CTRL+C signal
658    /// # Ok(())
659    /// # }
660    /// ```
661    pub fn kill_with(&mut self, signal: &str) -> Result<()> {
662        self.cancelled = true;
663        utils::trace_lazy("ProcessRunner", || format!("kill | signal={signal}"));
664
665        let Some(child) = self.child.as_mut() else {
666            return Ok(());
667        };
668
669        // Windows has no signals to deliver and no handler for the child to
670        // run, so there is nothing to grant a grace period to: the forceful
671        // stop is the only way to end the process.
672        // The `#[cfg(unix)]` block below is stripped on Windows, which leaves
673        // this one as the function's tail expression - hence no `return`.
674        #[cfg(not(unix))]
675        {
676            let _ = signal;
677            child.start_kill()?;
678            Ok(())
679        }
680
681        // Without a pid the process never spawned (or was already reaped);
682        // fall back to the forceful stop so `kill()` still terminates it.
683        #[cfg(unix)]
684        {
685            let Some(pid) = child.id() else {
686                child.start_kill()?;
687                return Ok(());
688            };
689
690            // `SIGKILL` cannot be handled, so there is nothing to wait for.
691            //
692            // A zero grace period means the child is given no opportunity to
693            // handle the signal either, so the requested signal is not
694            // delivered at all. Anything done between it and `SIGKILL` - even a
695            // single syscall - is a window the child can be scheduled in, which
696            // made "no grace" a race the child occasionally won rather than a
697            // guarantee. The reported exit code still comes from the signal
698            // that was requested.
699            let grace = self.options.kill_grace_ms;
700            let delivery = if self.own_process_group {
701                signal::Delivery::ProcessAndGroup
702            } else {
703                signal::Delivery::ProcessOnly
704            };
705            if grace == 0 || signal == "SIGKILL" {
706                signal::send_signal_to_process(pid, "SIGKILL", delivery);
707                let _ = child.start_kill();
708                return Ok(());
709            }
710
711            signal::send_signal_to_process(pid, signal, delivery);
712
713            // Otherwise escalate in the background so the child keeps its grace
714            // period without blocking the caller, which may not be inside an
715            // await point.
716            tokio::spawn(async move {
717                tokio::time::sleep(std::time::Duration::from_millis(grace)).await;
718                // Best effort: if the child already exited on the first signal
719                // this delivery simply fails, and the pid has not been reused
720                // because the `Child` handle above has not reaped it yet. That
721                // unreaped leader is also what keeps the group id alive, so the
722                // group delivery still reaches a grandchild that outlived it.
723                signal::send_signal_to_process(pid, "SIGKILL", delivery);
724            });
725
726            Ok(())
727        }
728    }
729
730    /// Check if the process is finished
731    pub fn is_finished(&self) -> bool {
732        self.finished
733    }
734
735    /// Get the result if available
736    pub fn result(&self) -> Option<&CommandResult> {
737        self.result.as_ref()
738    }
739
740    /// Process id of the command, or `None` when there is no operating system
741    /// process to identify.
742    ///
743    /// It is `None` before the command starts, and stays `None` for built-in
744    /// (virtual) commands such as `echo` or `sleep`, which run inside this
745    /// process and never spawn a child. Once a real command has been spawned
746    /// the value is stable: it remains readable after the command finishes,
747    /// unlike the child handle, which [`run`](Self::run) consumes.
748    ///
749    /// Mirrors the JavaScript `runner.pid` property.
750    ///
751    /// ```no_run
752    /// use command_stream::{ProcessRunner, RunOptions};
753    ///
754    /// # #[tokio::main]
755    /// # async fn main() -> command_stream::Result<()> {
756    /// let mut runner = ProcessRunner::new("/bin/sleep 1", RunOptions::default());
757    /// runner.start().await?;
758    /// println!("running as pid {:?}", runner.pid());
759    /// runner.run().await?;
760    /// println!("still readable: {:?}", runner.pid());
761    /// # Ok(())
762    /// # }
763    /// ```
764    pub fn pid(&self) -> Option<u32> {
765        self.pid
766    }
767
768    /// Get the command string
769    pub fn command(&self) -> &str {
770        &self.command
771    }
772
773    /// Get the options
774    pub fn options(&self) -> &RunOptions {
775        &self.options
776    }
777}
778
779/// Execute a command and return the result
780///
781/// This is the main entry point for simple command execution.
782/// Named `run` instead of `$` since `$` is not a valid Rust identifier.
783pub async fn run(command: impl Into<String>) -> Result<CommandResult> {
784    let mut runner = ProcessRunner::new(command, RunOptions::default());
785    runner.run().await
786}
787
788/// Alias for `run` function - for JavaScript-like API feel
789/// Since `$` is not valid in Rust, this provides a similar short name
790pub use run as execute;
791
792/// Execute a command with custom options
793pub async fn exec(command: impl Into<String>, options: RunOptions) -> Result<CommandResult> {
794    let mut runner = ProcessRunner::new(command, options);
795    runner.run().await
796}
797
798/// Create a new process runner without starting it
799pub fn create(command: impl Into<String>, options: RunOptions) -> ProcessRunner {
800    ProcessRunner::new(command, options)
801}
802
803/// Execute a command synchronously (blocking)
804pub fn run_sync(command: impl Into<String>) -> Result<CommandResult> {
805    let rt = tokio::runtime::Runtime::new()?;
806    rt.block_on(run(command))
807}
808
809// Tests are located in tests/ directory for better organization