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