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