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