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