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