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::{AsyncBufReadExt, AsyncWriteExt, BufReader};
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
108fn fallback_cwd() -> PathBuf {
109    std::env::var_os("HOME")
110        .or_else(|| std::env::var_os("USERPROFILE"))
111        .map(PathBuf::from)
112        .filter(|path| path.is_dir())
113        .unwrap_or_else(std::env::temp_dir)
114}
115
116/// Resolve a working directory that is safe to spawn a child process in.
117///
118/// When no explicit cwd is requested the child normally inherits the parent's
119/// working directory. But if that directory has been deleted or become
120/// inaccessible (the "getcwd() failed" scenario from issue #44), inheriting it
121/// makes the OS-level spawn fail. In that case fall back to a directory that is
122/// known to exist so the command still runs.
123///
124/// Normal behavior is preserved: when an explicit cwd is given, or when the
125/// inherited working directory is valid, this returns the requested value
126/// (`None` meaning "inherit").
127fn resolve_spawn_cwd(cwd: Option<&PathBuf>) -> Option<PathBuf> {
128    // An explicit directory is always honored as-is.
129    if let Some(c) = cwd {
130        return Some(c.clone());
131    }
132
133    // No explicit cwd: we would inherit the parent's working directory. Make
134    // sure that directory is actually usable before relying on inheritance.
135    match std::env::current_dir() {
136        Ok(_) => None,
137        Err(e) => {
138            let fallback = fallback_cwd();
139            trace(
140                "ProcessRunner",
141                &format!(
142                    "current_dir() failed ({}); spawning in fallback directory {}",
143                    e,
144                    fallback.display()
145                ),
146            );
147            Some(fallback)
148        }
149    }
150}
151
152/// Error type for command-stream operations
153#[derive(Debug, thiserror::Error)]
154pub enum Error {
155    #[error("IO error: {0}")]
156    Io(#[from] std::io::Error),
157
158    #[error("Command failed with exit code {code}: {message}")]
159    CommandFailed { code: i32, message: String },
160
161    #[error("Command not found: {0}")]
162    CommandNotFound(String),
163
164    #[error("Parse error: {0}")]
165    ParseError(String),
166
167    #[error("Cancelled")]
168    Cancelled,
169}
170
171/// Result type for command-stream operations
172pub type Result<T> = std::result::Result<T, Error>;
173
174/// Options for command execution
175#[derive(Debug, Clone)]
176pub struct RunOptions {
177    /// Mirror output to parent stdout/stderr
178    pub mirror: bool,
179    /// Capture output in result
180    pub capture: bool,
181    /// Standard input handling
182    pub stdin: StdinOption,
183    /// Working directory
184    pub cwd: Option<PathBuf>,
185    /// Environment variables
186    pub env: Option<HashMap<String, String>>,
187    /// Interactive mode (TTY forwarding)
188    pub interactive: bool,
189    /// Enable shell operator parsing
190    pub shell_operators: bool,
191    /// Enable tracing for this command
192    pub trace: bool,
193}
194
195impl Default for RunOptions {
196    fn default() -> Self {
197        RunOptions {
198            mirror: true,
199            capture: true,
200            stdin: StdinOption::Inherit,
201            cwd: None,
202            env: None,
203            interactive: false,
204            shell_operators: true,
205            trace: true,
206        }
207    }
208}
209
210/// Standard input options
211#[derive(Debug, Clone)]
212pub enum StdinOption {
213    /// Inherit from parent process
214    Inherit,
215    /// Pipe (allow writing to stdin)
216    Pipe,
217    /// Provide string content
218    Content(String),
219    /// Null device
220    Null,
221}
222
223/// A running or completed process
224pub struct ProcessRunner {
225    command: String,
226    options: RunOptions,
227    child: Option<Child>,
228    result: Option<CommandResult>,
229    started: bool,
230    finished: bool,
231    cancelled: bool,
232    output_tx: Option<mpsc::Sender<StreamChunk>>,
233    // Held, never read: dropping the receiver would close the channel, and
234    // streaming virtual commands treat a closed channel as "stop now" (see
235    // `commands::yes`, which loops until `output_tx.send` fails). Keeping it
236    // alive is what gives those commands their run-until-cancelled behaviour.
237    #[allow(dead_code)]
238    output_rx: Option<mpsc::Receiver<StreamChunk>>,
239}
240
241impl ProcessRunner {
242    /// Create a new process runner
243    pub fn new(command: impl Into<String>, options: RunOptions) -> Self {
244        let (tx, rx) = mpsc::channel(1024);
245        ProcessRunner {
246            command: command.into(),
247            options,
248            child: None,
249            result: None,
250            started: false,
251            finished: false,
252            cancelled: false,
253            output_tx: Some(tx),
254            output_rx: Some(rx),
255        }
256    }
257
258    /// Start the process
259    pub async fn start(&mut self) -> Result<()> {
260        if self.started {
261            return Ok(());
262        }
263        self.started = true;
264
265        utils::trace_lazy("ProcessRunner", || {
266            format!("Starting command: {}", self.command)
267        });
268
269        // Check if this is a virtual command. Backslash escapes are removed by a
270        // real shell but not by the whitespace splitting used for virtual
271        // command args, so such commands always go to the system shell (#49).
272        // The same applies to redirection and expansions: whitespace splitting
273        // would hand `>`, `out.txt` to the virtual command as two ordinary
274        // arguments, so `echo hello > out.txt` printed the redirection instead
275        // of writing the file, and `git push ... 2>&1` reported success while
276        // nothing was pushed (#46).
277        let first_word = if has_shell_escapes(&self.command) || needs_real_shell(&self.command) {
278            ""
279        } else {
280            self.command.split_whitespace().next().unwrap_or("")
281        };
282        if let Some(result) = self.try_virtual_command(first_word).await {
283            self.result = Some(result);
284            self.finished = true;
285            return Ok(());
286        }
287        // Parse command for shell operators (for future use with virtual command pipelines)
288        let _parsed = if self.options.shell_operators && !needs_real_shell(&self.command) {
289            parse_shell_command(&self.command)
290        } else {
291            None
292        };
293
294        // Execute via real shell if needed
295        let shell = find_available_shell();
296
297        let mut cmd = Command::new(&shell.cmd);
298        for arg in &shell.args {
299            cmd.arg(arg);
300        }
301        cmd.arg(utils::with_exported_process_context(
302            &self.command,
303            self.options.env.as_ref(),
304        ));
305
306        // Configure stdin
307        match &self.options.stdin {
308            StdinOption::Inherit => {
309                cmd.stdin(Stdio::inherit());
310            }
311            StdinOption::Pipe => {
312                cmd.stdin(Stdio::piped());
313            }
314            StdinOption::Content(_) => {
315                cmd.stdin(Stdio::piped());
316            }
317            StdinOption::Null => {
318                cmd.stdin(Stdio::null());
319            }
320        }
321
322        // Configure stdout/stderr
323        if self.options.capture || self.options.mirror {
324            cmd.stdout(Stdio::piped());
325            cmd.stderr(Stdio::piped());
326        } else {
327            cmd.stdout(Stdio::inherit());
328            cmd.stderr(Stdio::inherit());
329        }
330
331        // Set working directory. Fall back to a valid directory when the
332        // inherited working directory has been deleted (issue #44).
333        if let Some(cwd) = resolve_spawn_cwd(self.options.cwd.as_ref()) {
334            cmd.current_dir(cwd);
335        }
336
337        // Set environment
338        if let Some(ref env_vars) = self.options.env {
339            for (key, value) in env_vars {
340                cmd.env(key, value);
341            }
342        }
343
344        // Spawn the process
345        let child = cmd.spawn()?;
346        self.child = Some(child);
347
348        Ok(())
349    }
350
351    /// Run the process to completion
352    pub async fn run(&mut self) -> Result<CommandResult> {
353        self.start().await?;
354
355        if let Some(result) = &self.result {
356            return Ok(result.clone());
357        }
358
359        let mut child = self
360            .child
361            .take()
362            .ok_or_else(|| Error::Io(std::io::Error::other("Process not started")))?;
363
364        // Handle stdin content if provided
365        if let StdinOption::Content(ref content) = self.options.stdin {
366            if let Some(mut stdin) = child.stdin.take() {
367                let content = content.clone();
368                tokio::spawn(async move {
369                    let _ = stdin.write_all(content.as_bytes()).await;
370                    let _ = stdin.shutdown().await;
371                });
372            }
373        }
374
375        // Collect output
376        let mut stdout_content = String::new();
377        let mut stderr_content = String::new();
378
379        if let Some(stdout) = child.stdout.take() {
380            let mut reader = BufReader::new(stdout).lines();
381            while let Ok(Some(line)) = reader.next_line().await {
382                if self.options.mirror {
383                    println!("{}", line);
384                }
385                stdout_content.push_str(&line);
386                stdout_content.push('\n');
387            }
388        }
389
390        if let Some(stderr) = child.stderr.take() {
391            let mut reader = BufReader::new(stderr).lines();
392            while let Ok(Some(line)) = reader.next_line().await {
393                if self.options.mirror {
394                    eprintln!("{}", line);
395                }
396                stderr_content.push_str(&line);
397                stderr_content.push('\n');
398            }
399        }
400
401        let status = child.wait().await?;
402        let code = status.code().unwrap_or(-1);
403
404        let result = CommandResult {
405            stdout: stdout_content,
406            stderr: stderr_content,
407            code,
408        };
409
410        self.result = Some(result.clone());
411        self.finished = true;
412
413        Ok(result)
414    }
415
416    /// Try to execute as a virtual command
417    async fn try_virtual_command(&self, cmd_name: &str) -> Option<CommandResult> {
418        if !commands::are_virtual_commands_enabled() {
419            return None;
420        }
421
422        // An empty command name means the caller already decided this command
423        // must go to a real shell (redirection, expansions, escapes). Bail out
424        // before tokenizing so we neither waste work nor parse shell syntax we
425        // deliberately delegate.
426        if cmd_name.is_empty() {
427            return None;
428        }
429
430        // Parse args from command string, respecting quotes and performing
431        // POSIX quote removal so `echo label:'help wanted'` reaches the built-in
432        // as the single argument `label:help wanted` (issue #48).
433        let words = shell_parser::split_command_words(&self.command);
434        let args: Vec<String> = words.into_iter().skip(1).collect();
435
436        let ctx = CommandContext {
437            args,
438            stdin: match &self.options.stdin {
439                StdinOption::Content(s) => Some(s.clone()),
440                _ => None,
441            },
442            cwd: self.options.cwd.clone(),
443            env: self.options.env.clone(),
444            output_tx: self.output_tx.clone(),
445            is_cancelled: None,
446        };
447
448        match cmd_name {
449            "echo" => Some(commands::echo(ctx).await),
450            "pwd" => Some(commands::pwd(ctx).await),
451            "cd" => Some(commands::cd::resolve_cd(ctx).await.0),
452            "true" => Some(commands::r#true(ctx).await),
453            "false" => Some(commands::r#false(ctx).await),
454            "sleep" => Some(commands::sleep(ctx).await),
455            "cat" => Some(commands::cat(ctx).await),
456            "ls" => Some(commands::ls(ctx).await),
457            "mkdir" => Some(commands::mkdir(ctx).await),
458            "rm" => Some(commands::rm(ctx).await),
459            "touch" => Some(commands::touch(ctx).await),
460            "cp" => Some(commands::cp(ctx).await),
461            "mv" => Some(commands::mv(ctx).await),
462            "basename" => Some(commands::basename(ctx).await),
463            "dirname" => Some(commands::dirname(ctx).await),
464            "env" => Some(commands::env(ctx).await),
465            "exit" => Some(commands::exit(ctx).await),
466            "which" => Some(commands::which(ctx).await),
467            "yes" => Some(commands::yes(ctx).await),
468            "seq" => Some(commands::seq(ctx).await),
469            "test" => Some(commands::test(ctx).await),
470            _ => None,
471        }
472    }
473
474    /// Kill the process
475    pub fn kill(&mut self) -> Result<()> {
476        self.cancelled = true;
477        if let Some(ref mut child) = self.child {
478            child.start_kill()?;
479        }
480        Ok(())
481    }
482
483    /// Check if the process is finished
484    pub fn is_finished(&self) -> bool {
485        self.finished
486    }
487
488    /// Get the result if available
489    pub fn result(&self) -> Option<&CommandResult> {
490        self.result.as_ref()
491    }
492
493    /// Get the command string
494    pub fn command(&self) -> &str {
495        &self.command
496    }
497
498    /// Get the options
499    pub fn options(&self) -> &RunOptions {
500        &self.options
501    }
502}
503
504/// Shell configuration
505#[derive(Debug, Clone)]
506struct ShellConfig {
507    cmd: String,
508    args: Vec<String>,
509}
510
511/// Find an available shell
512fn find_available_shell() -> ShellConfig {
513    let is_windows = cfg!(windows);
514
515    if is_windows {
516        // Windows shells
517        let shells = [
518            ("cmd.exe", vec!["/c"]),
519            ("powershell.exe", vec!["-Command"]),
520        ];
521
522        for (cmd, args) in shells {
523            if which::which(cmd).is_ok() {
524                return ShellConfig {
525                    cmd: cmd.to_string(),
526                    args: args.into_iter().map(String::from).collect(),
527                };
528            }
529        }
530
531        ShellConfig {
532            cmd: "cmd.exe".to_string(),
533            args: vec!["/c".to_string()],
534        }
535    } else {
536        // Unix shells
537        let shells = [
538            ("/bin/sh", vec!["-c"]),
539            ("/usr/bin/sh", vec!["-c"]),
540            ("/bin/bash", vec!["-c"]),
541            ("sh", vec!["-c"]),
542        ];
543
544        for (cmd, args) in shells {
545            if std::path::Path::new(cmd).exists() || which::which(cmd).is_ok() {
546                return ShellConfig {
547                    cmd: cmd.to_string(),
548                    args: args.into_iter().map(String::from).collect(),
549                };
550            }
551        }
552
553        ShellConfig {
554            cmd: "/bin/sh".to_string(),
555            args: vec!["-c".to_string()],
556        }
557    }
558}
559
560/// Execute a command and return the result
561///
562/// This is the main entry point for simple command execution.
563/// Named `run` instead of `$` since `$` is not a valid Rust identifier.
564pub async fn run(command: impl Into<String>) -> Result<CommandResult> {
565    let mut runner = ProcessRunner::new(command, RunOptions::default());
566    runner.run().await
567}
568
569/// Alias for `run` function - for JavaScript-like API feel
570/// Since `$` is not valid in Rust, this provides a similar short name
571pub use run as execute;
572
573/// Execute a command with custom options
574pub async fn exec(command: impl Into<String>, options: RunOptions) -> Result<CommandResult> {
575    let mut runner = ProcessRunner::new(command, options);
576    runner.run().await
577}
578
579/// Create a new process runner without starting it
580pub fn create(command: impl Into<String>, options: RunOptions) -> ProcessRunner {
581    ProcessRunner::new(command, options)
582}
583
584/// Execute a command synchronously (blocking)
585pub fn run_sync(command: impl Into<String>) -> Result<CommandResult> {
586    let rt = tokio::runtime::Runtime::new()?;
587    rt.block_on(run(command))
588}
589
590// Tests are located in tests/ directory for better organization