Skip to main content

command_stream/
lib.rs

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