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