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