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