Skip to main content

command_stream/
lib.rs

1//! # command-stream
2//!
3//! Modern shell command execution library with streaming, async iteration, and event support.
4//!
5//! This library provides a Rust equivalent to the JavaScript command-stream library,
6//! offering powerful shell command execution with streaming capabilities.
7//!
8//! ## Features
9//!
10//! - Async command execution with tokio
11//! - Streaming output via async iterators
12//! - Event-based output handling (on, once, emit)
13//! - Virtual commands for common operations (cat, ls, mkdir, etc.)
14//! - Shell operator support (&&, ||, ;, |)
15//! - Pipeline support with `.pipe()` method and `Pipeline` builder
16//! - Global state management for shell settings
17//! - `cmd!` macro for ergonomic command creation (similar to JS `$` tagged template literals)
18//! - Cross-platform support
19//!
20//! ## Module Organization
21//!
22//! The codebase follows a modular architecture similar to the JavaScript implementation:
23//!
24//! - `ansi` - ANSI escape code handling utilities
25//! - `commands` - Virtual command implementations
26//! - `events` - Event emitter for stream events
27//! - `macros` - The `cmd!` macro for ergonomic command creation
28//! - `pipeline` - Pipeline execution support
29//! - `quote` - Shell quoting utilities
30//! - `shell_parser` - Shell command parsing
31//! - `state` - Global state management
32//! - `stream` - Async streaming and iteration support
33//! - `trace` - Logging and tracing utilities
34//! - `utils` - Command results and virtual command helpers
35//!
36//! ## Quick Start
37//!
38//! ```rust,no_run
39//! use command_stream::{run, cmd};
40//!
41//! #[tokio::main]
42//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
43//!     // Execute a simple command
44//!     let result = run("echo hello world").await?;
45//!     println!("{}", result.stdout);
46//!
47//!     // Using the cmd! macro (similar to JS $ tagged template)
48//!     let name = "world";
49//!     let result = cmd!("echo hello {}", name).await?;
50//!     println!("{}", result.stdout);
51//!
52//!     // Using pipelines
53//!     use command_stream::Pipeline;
54//!     let result = Pipeline::new()
55//!         .add("echo hello world")
56//!         .add("grep world")
57//!         .run()
58//!         .await?;
59//!
60//!     Ok(())
61//! }
62//! ```
63
64// Modular utility modules (following JavaScript modular pattern)
65pub mod ansi;
66pub mod events;
67#[doc(hidden)]
68pub mod macros;
69pub mod pipeline;
70pub mod quote;
71pub mod state;
72pub mod stream;
73pub mod terminal;
74pub mod trace;
75
76// Core modules
77pub mod commands;
78pub mod shell_parser;
79pub mod utils;
80
81use std::collections::HashMap;
82use std::path::PathBuf;
83use std::process::Stdio;
84use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt};
85use tokio::process::{Child, Command};
86use tokio::sync::mpsc;
87
88pub use commands::{CommandContext, StreamChunk};
89pub use shell_parser::{needs_real_shell, parse_shell_command, ParsedCommand};
90pub use utils::{CommandResult, VirtualUtils};
91
92// Re-export modular utilities at crate root for convenient access
93pub use ansi::{AnsiConfig, AnsiUtils};
94pub use events::{EventData, EventType, StreamEmitter};
95pub use pipeline::{Pipeline, PipelineBuilder, PipelineExt};
96pub use quote::{
97    escape_for_double_quotes, escape_for_single_quotes, has_shell_escapes,
98    is_pre_quoted_passthrough_enabled, is_quote_context_enabled, quote, quote_for_context,
99    scan_quote_context, QuoteContext,
100};
101pub use state::{
102    get_shell_settings, global_state, reset_global_state, set_shell_option, unset_shell_option,
103    GlobalState, ShellSettings,
104};
105pub use stream::{AsyncIterator, IntoStream, OutputChunk, OutputStream, StreamingRunner};
106pub use trace::trace;
107
108#[derive(Clone, Copy)]
109enum ChildOutput {
110    Stdout,
111    Stderr,
112}
113
114/// Read child output as byte chunks so capture does not invent a trailing newline.
115///
116/// stdout and stderr use separate futures in `ProcessRunner::run`, preventing
117/// either pipe from filling while the other is being drained. Mirroring keeps
118/// the original bytes too, including output that does not end in a newline.
119async fn collect_child_output<R>(
120    reader: Option<R>,
121    mirror: bool,
122    target: ChildOutput,
123) -> std::io::Result<Vec<u8>>
124where
125    R: AsyncRead + Unpin,
126{
127    let Some(mut reader) = reader else {
128        return Ok(Vec::new());
129    };
130    let mut collected = Vec::new();
131    let mut buffer = [0_u8; 8192];
132
133    loop {
134        let count = reader.read(&mut buffer).await?;
135        if count == 0 {
136            break;
137        }
138
139        let chunk = &buffer[..count];
140        collected.extend_from_slice(chunk);
141        if mirror {
142            match target {
143                ChildOutput::Stdout => {
144                    let mut output = std::io::stdout().lock();
145                    let _ = std::io::Write::write_all(&mut output, chunk);
146                    let _ = std::io::Write::flush(&mut output);
147                }
148                ChildOutput::Stderr => {
149                    let mut output = std::io::stderr().lock();
150                    let _ = std::io::Write::write_all(&mut output, chunk);
151                    let _ = std::io::Write::flush(&mut output);
152                }
153            }
154        }
155    }
156
157    Ok(collected)
158}
159
160fn fallback_cwd() -> PathBuf {
161    std::env::var_os("HOME")
162        .or_else(|| std::env::var_os("USERPROFILE"))
163        .map(PathBuf::from)
164        .filter(|path| path.is_dir())
165        .unwrap_or_else(std::env::temp_dir)
166}
167
168/// Resolve a working directory that is safe to spawn a child process in.
169///
170/// When no explicit cwd is requested the child normally inherits the parent's
171/// working directory. But if that directory has been deleted or become
172/// inaccessible (the "getcwd() failed" scenario from issue #44), inheriting it
173/// makes the OS-level spawn fail. In that case fall back to a directory that is
174/// known to exist so the command still runs.
175///
176/// Normal behavior is preserved: when an explicit cwd is given, or when the
177/// inherited working directory is valid, this returns the requested value
178/// (`None` meaning "inherit").
179fn resolve_spawn_cwd(cwd: Option<&PathBuf>) -> Option<PathBuf> {
180    // An explicit directory is always honored as-is.
181    if let Some(c) = cwd {
182        return Some(c.clone());
183    }
184
185    // No explicit cwd: we would inherit the parent's working directory. Make
186    // sure that directory is actually usable before relying on inheritance.
187    match std::env::current_dir() {
188        Ok(_) => None,
189        Err(e) => {
190            let fallback = fallback_cwd();
191            trace(
192                "ProcessRunner",
193                &format!(
194                    "current_dir() failed ({}); spawning in fallback directory {}",
195                    e,
196                    fallback.display()
197                ),
198            );
199            Some(fallback)
200        }
201    }
202}
203
204/// Error type for command-stream operations
205#[derive(Debug, thiserror::Error)]
206pub enum Error {
207    #[error("IO error: {0}")]
208    Io(#[from] std::io::Error),
209
210    #[error("Command failed with exit code {code}: {message}")]
211    CommandFailed { code: i32, message: String },
212
213    #[error("Command not found: {0}")]
214    CommandNotFound(String),
215
216    #[error("Parse error: {0}")]
217    ParseError(String),
218
219    #[error("Cancelled")]
220    Cancelled,
221}
222
223impl Error {
224    /// Build a [`Error::CommandFailed`] for a command that exited with `code`.
225    pub fn command_failed(code: i32, message: impl Into<String>) -> Self {
226        Error::CommandFailed {
227            code,
228            message: message.into(),
229        }
230    }
231
232    /// Exit status carried by the error, when the failure has one.
233    ///
234    /// Mirrors the `error.code` property of the JavaScript implementation
235    /// (issue #38). Failures that never reached a child process, such as parse
236    /// errors, report `None`.
237    pub fn code(&self) -> Option<i32> {
238        match self {
239            Error::CommandFailed { code, .. } => Some(*code),
240            // `command not found` is 127 in POSIX shells, which is also what
241            // the JavaScript implementation reports for a missing executable.
242            Error::CommandNotFound(_) => Some(127),
243            Error::Io(error) => match error.kind() {
244                std::io::ErrorKind::NotFound => Some(127),
245                std::io::ErrorKind::PermissionDenied => Some(126),
246                _ => None,
247            },
248            // A cancelled command is terminated with SIGINT (128 + 2).
249            Error::Cancelled => Some(130),
250            Error::ParseError(_) => None,
251        }
252    }
253
254    /// Alias for [`code`](Self::code).
255    ///
256    /// Node.js `child_process` names this property `code`, while execa, zx,
257    /// nano-spawn, and Bun Shell name it `exitCode`. command-stream exposes
258    /// both spellings in every language (issue #38).
259    pub fn exit_code(&self) -> Option<i32> {
260        self.code()
261    }
262}
263
264/// Result type for command-stream operations
265pub type Result<T> = std::result::Result<T, Error>;
266
267/// Options for command execution
268#[derive(Debug, Clone)]
269pub struct RunOptions {
270    /// Mirror output to parent stdout/stderr
271    pub mirror: bool,
272    /// Capture output in result
273    pub capture: bool,
274    /// Standard input handling
275    pub stdin: StdinOption,
276    /// Working directory
277    pub cwd: Option<PathBuf>,
278    /// Environment variables
279    pub env: Option<HashMap<String, String>>,
280    /// Interactive mode (TTY forwarding)
281    pub interactive: bool,
282    /// Enable shell operator parsing
283    pub shell_operators: bool,
284    /// Enable tracing for this command
285    pub trace: bool,
286}
287
288impl Default for RunOptions {
289    fn default() -> Self {
290        RunOptions {
291            mirror: true,
292            capture: true,
293            stdin: StdinOption::Inherit,
294            cwd: None,
295            env: None,
296            interactive: false,
297            shell_operators: true,
298            trace: true,
299        }
300    }
301}
302
303/// Standard input options
304#[derive(Debug, Clone)]
305pub enum StdinOption {
306    /// Inherit from parent process
307    Inherit,
308    /// Pipe (allow writing to stdin)
309    Pipe,
310    /// Provide string content
311    Content(String),
312    /// Null device
313    Null,
314}
315
316/// A running or completed process
317pub struct ProcessRunner {
318    command: String,
319    options: RunOptions,
320    child: Option<Child>,
321    result: Option<CommandResult>,
322    started: bool,
323    finished: bool,
324    cancelled: bool,
325    output_tx: Option<mpsc::Sender<StreamChunk>>,
326    // Held, never read: dropping the receiver would close the channel, and
327    // streaming virtual commands treat a closed channel as "stop now" (see
328    // `commands::yes`, which loops until `output_tx.send` fails). Keeping it
329    // alive is what gives those commands their run-until-cancelled behaviour.
330    #[allow(dead_code)]
331    output_rx: Option<mpsc::Receiver<StreamChunk>>,
332}
333
334impl ProcessRunner {
335    /// Create a new process runner
336    pub fn new(command: impl Into<String>, options: RunOptions) -> Self {
337        let (tx, rx) = mpsc::channel(1024);
338        ProcessRunner {
339            command: command.into(),
340            options,
341            child: None,
342            result: None,
343            started: false,
344            finished: false,
345            cancelled: false,
346            output_tx: Some(tx),
347            output_rx: Some(rx),
348        }
349    }
350
351    /// Start the process
352    pub async fn start(&mut self) -> Result<()> {
353        if self.started {
354            return Ok(());
355        }
356        self.started = true;
357
358        utils::trace_lazy("ProcessRunner", || {
359            format!("Starting command: {}", self.command)
360        });
361
362        // Check if this is a virtual command. Backslash escapes are removed by a
363        // real shell but not by the whitespace splitting used for virtual
364        // command args, so such commands always go to the system shell (#49).
365        // The same applies to redirection and expansions: whitespace splitting
366        // would hand `>`, `out.txt` to the virtual command as two ordinary
367        // arguments, so `echo hello > out.txt` printed the redirection instead
368        // of writing the file, and `git push ... 2>&1` reported success while
369        // nothing was pushed (#46).
370        let first_word = if has_shell_escapes(&self.command) || needs_real_shell(&self.command) {
371            ""
372        } else {
373            self.command.split_whitespace().next().unwrap_or("")
374        };
375        if let Some(result) = self.try_virtual_command(first_word).await {
376            self.result = Some(result);
377            self.finished = true;
378            return Ok(());
379        }
380        // Parse command for shell operators (for future use with virtual command pipelines)
381        let _parsed = if self.options.shell_operators && !needs_real_shell(&self.command) {
382            parse_shell_command(&self.command)
383        } else {
384            None
385        };
386
387        // Execute via real shell if needed
388        let shell = find_available_shell();
389
390        let mut cmd = Command::new(&shell.cmd);
391        for arg in &shell.args {
392            cmd.arg(arg);
393        }
394        utils::append_shell_command(&mut cmd, &self.command, self.options.env.as_ref());
395
396        // Configure stdin
397        match &self.options.stdin {
398            StdinOption::Inherit => {
399                cmd.stdin(Stdio::inherit());
400            }
401            StdinOption::Pipe => {
402                cmd.stdin(Stdio::piped());
403            }
404            StdinOption::Content(_) => {
405                cmd.stdin(Stdio::piped());
406            }
407            StdinOption::Null => {
408                cmd.stdin(Stdio::null());
409            }
410        }
411
412        // Configure stdout/stderr
413        if self.options.capture || self.options.mirror {
414            cmd.stdout(Stdio::piped());
415            cmd.stderr(Stdio::piped());
416        } else {
417            cmd.stdout(Stdio::inherit());
418            cmd.stderr(Stdio::inherit());
419        }
420
421        // Set working directory. Fall back to a valid directory when the
422        // inherited working directory has been deleted (issue #44).
423        if let Some(cwd) = resolve_spawn_cwd(self.options.cwd.as_ref()) {
424            cmd.current_dir(cwd);
425        }
426
427        // Set environment
428        if let Some(ref env_vars) = self.options.env {
429            for (key, value) in env_vars {
430                cmd.env(key, value);
431            }
432        }
433
434        // Spawn the process
435        let child = cmd.spawn()?;
436        self.child = Some(child);
437
438        Ok(())
439    }
440
441    /// Run the process to completion
442    pub async fn run(&mut self) -> Result<CommandResult> {
443        self.start().await?;
444
445        if let Some(result) = &self.result {
446            return Ok(result.clone());
447        }
448
449        let mut child = self
450            .child
451            .take()
452            .ok_or_else(|| Error::Io(std::io::Error::other("Process not started")))?;
453
454        // Handle stdin content if provided
455        if let StdinOption::Content(ref content) = self.options.stdin {
456            if let Some(mut stdin) = child.stdin.take() {
457                let content = content.clone();
458                tokio::spawn(async move {
459                    let _ = stdin.write_all(content.as_bytes()).await;
460                    let _ = stdin.shutdown().await;
461                });
462            }
463        }
464
465        // Drain both pipes concurrently and preserve their newline framing. The
466        // previous line reader appended `\n` to every final line, changing
467        // output from commands such as `printf` that omit a newline (issue #37).
468        let stdout = child.stdout.take();
469        let stderr = child.stderr.take();
470        let collected = tokio::try_join!(
471            collect_child_output(stdout, self.options.mirror, ChildOutput::Stdout),
472            collect_child_output(stderr, self.options.mirror, ChildOutput::Stderr),
473        );
474        let (stdout, stderr) = match collected {
475            Ok(output) => output,
476            Err(error) => {
477                // `try_join!` drops the other pipe reader after an error. Stop
478                // and reap the child so it cannot remain blocked on that pipe.
479                let _ = child.start_kill();
480                let _ = child.wait().await;
481                return Err(error.into());
482            }
483        };
484
485        let status = child.wait().await?;
486        let code = status.code().unwrap_or(-1);
487
488        let result = CommandResult {
489            stdout: String::from_utf8_lossy(&stdout).into_owned(),
490            stderr: String::from_utf8_lossy(&stderr).into_owned(),
491            code,
492        };
493
494        self.result = Some(result.clone());
495        self.finished = true;
496
497        Ok(result)
498    }
499
500    /// Try to execute as a virtual command
501    async fn try_virtual_command(&self, cmd_name: &str) -> Option<CommandResult> {
502        if !commands::are_virtual_commands_enabled() {
503            return None;
504        }
505
506        // An empty command name means the caller already decided this command
507        // must go to a real shell (redirection, expansions, escapes). Bail out
508        // before tokenizing so we neither waste work nor parse shell syntax we
509        // deliberately delegate.
510        if cmd_name.is_empty() {
511            return None;
512        }
513
514        // Parse args from command string, respecting quotes and performing
515        // POSIX quote removal so `echo label:'help wanted'` reaches the built-in
516        // as the single argument `label:help wanted` (issue #48).
517        let words = shell_parser::split_command_words(&self.command);
518        let args: Vec<String> = words.into_iter().skip(1).collect();
519
520        let ctx = CommandContext {
521            args,
522            stdin: match &self.options.stdin {
523                StdinOption::Content(s) => Some(s.clone()),
524                _ => None,
525            },
526            cwd: self.options.cwd.clone(),
527            env: self.options.env.clone(),
528            output_tx: self.output_tx.clone(),
529            is_cancelled: None,
530        };
531
532        match cmd_name {
533            "echo" => Some(commands::echo(ctx).await),
534            "pwd" => Some(commands::pwd(ctx).await),
535            "cd" => Some(commands::cd::resolve_cd(ctx).await.0),
536            "true" => Some(commands::r#true(ctx).await),
537            "false" => Some(commands::r#false(ctx).await),
538            "sleep" => Some(commands::sleep(ctx).await),
539            "cat" => Some(commands::cat(ctx).await),
540            "ls" => Some(commands::ls(ctx).await),
541            "mkdir" => Some(commands::mkdir(ctx).await),
542            "rm" => Some(commands::rm(ctx).await),
543            "touch" => Some(commands::touch(ctx).await),
544            "cp" => Some(commands::cp(ctx).await),
545            "mv" => Some(commands::mv(ctx).await),
546            "basename" => Some(commands::basename(ctx).await),
547            "dirname" => Some(commands::dirname(ctx).await),
548            "env" => Some(commands::env(ctx).await),
549            "exit" => Some(commands::exit(ctx).await),
550            "which" => Some(commands::which(ctx).await),
551            "yes" => Some(commands::yes(ctx).await),
552            "seq" => Some(commands::seq(ctx).await),
553            "tee" => Some(commands::tee(ctx).await),
554            "test" => Some(commands::test(ctx).await),
555            _ => None,
556        }
557    }
558
559    /// Kill the process
560    pub fn kill(&mut self) -> Result<()> {
561        self.cancelled = true;
562        if let Some(ref mut child) = self.child {
563            child.start_kill()?;
564        }
565        Ok(())
566    }
567
568    /// Check if the process is finished
569    pub fn is_finished(&self) -> bool {
570        self.finished
571    }
572
573    /// Get the result if available
574    pub fn result(&self) -> Option<&CommandResult> {
575        self.result.as_ref()
576    }
577
578    /// Get the command string
579    pub fn command(&self) -> &str {
580        &self.command
581    }
582
583    /// Get the options
584    pub fn options(&self) -> &RunOptions {
585        &self.options
586    }
587}
588
589/// Shell configuration
590#[derive(Debug, Clone)]
591struct ShellConfig {
592    cmd: String,
593    args: Vec<String>,
594}
595
596/// Find an available shell
597fn find_available_shell() -> ShellConfig {
598    let is_windows = cfg!(windows);
599
600    if is_windows {
601        // Windows shells
602        let shells = [
603            ("cmd.exe", vec!["/c"]),
604            ("powershell.exe", vec!["-Command"]),
605        ];
606
607        for (cmd, args) in shells {
608            if which::which(cmd).is_ok() {
609                return ShellConfig {
610                    cmd: cmd.to_string(),
611                    args: args.into_iter().map(String::from).collect(),
612                };
613            }
614        }
615
616        ShellConfig {
617            cmd: "cmd.exe".to_string(),
618            args: vec!["/c".to_string()],
619        }
620    } else {
621        // Unix shells
622        let shells = [
623            ("/bin/sh", vec!["-c"]),
624            ("/usr/bin/sh", vec!["-c"]),
625            ("/bin/bash", vec!["-c"]),
626            ("sh", vec!["-c"]),
627        ];
628
629        for (cmd, args) in shells {
630            if std::path::Path::new(cmd).exists() || which::which(cmd).is_ok() {
631                return ShellConfig {
632                    cmd: cmd.to_string(),
633                    args: args.into_iter().map(String::from).collect(),
634                };
635            }
636        }
637
638        ShellConfig {
639            cmd: "/bin/sh".to_string(),
640            args: vec!["-c".to_string()],
641        }
642    }
643}
644
645/// Execute a command and return the result
646///
647/// This is the main entry point for simple command execution.
648/// Named `run` instead of `$` since `$` is not a valid Rust identifier.
649pub async fn run(command: impl Into<String>) -> Result<CommandResult> {
650    let mut runner = ProcessRunner::new(command, RunOptions::default());
651    runner.run().await
652}
653
654/// Alias for `run` function - for JavaScript-like API feel
655/// Since `$` is not valid in Rust, this provides a similar short name
656pub use run as execute;
657
658/// Execute a command with custom options
659pub async fn exec(command: impl Into<String>, options: RunOptions) -> Result<CommandResult> {
660    let mut runner = ProcessRunner::new(command, options);
661    runner.run().await
662}
663
664/// Create a new process runner without starting it
665pub fn create(command: impl Into<String>, options: RunOptions) -> ProcessRunner {
666    ProcessRunner::new(command, options)
667}
668
669/// Execute a command synchronously (blocking)
670pub fn run_sync(command: impl Into<String>) -> Result<CommandResult> {
671    let rt = tokio::runtime::Runtime::new()?;
672    rt.block_on(run(command))
673}
674
675// Tests are located in tests/ directory for better organization