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) -> std::io::Result<Vec<u8>>
137where
138 R: AsyncRead + Unpin,
139{
140 let Some(mut reader) = reader else {
141 return Ok(Vec::new());
142 };
143 let mut collected = Vec::new();
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(collected)
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 = child.stdout.take();
536 let stderr = child.stderr.take();
537 let collected = tokio::try_join!(
538 collect_child_output(stdout, self.options.mirror, ChildOutput::Stdout),
539 collect_child_output(stderr, self.options.mirror, ChildOutput::Stderr),
540 );
541 let (stdout, stderr) = match collected {
542 Ok(output) => output,
543 Err(error) => {
544 // `try_join!` drops the other pipe reader after an error. Stop
545 // and reap the child so it cannot remain blocked on that pipe.
546 let _ = child.start_kill();
547 let _ = child.wait().await;
548 return Err(error.into());
549 }
550 };
551
552 let status = child.wait().await?;
553 let code = status.code().unwrap_or(-1);
554
555 let mut result = CommandResult::new(
556 String::from_utf8_lossy(&stdout).into_owned(),
557 String::from_utf8_lossy(&stderr).into_owned(),
558 code,
559 );
560 if let StdinOption::Content(ref content) = self.options.stdin {
561 result.stdin = crate::result_streams::CapturedInput::new(content.as_bytes().to_vec());
562 } else if !self.stdin_bytes.is_empty() {
563 result.stdin = crate::result_streams::CapturedInput::new(self.stdin_bytes.clone());
564 }
565
566 self.result = Some(result.clone());
567 self.finished = true;
568
569 Ok(result)
570 }
571
572 /// Try to execute as a virtual command
573 async fn try_virtual_command(&self, cmd_name: &str) -> Option<CommandResult> {
574 if !commands::are_virtual_commands_enabled() {
575 return None;
576 }
577
578 // An empty command name means the caller already decided this command
579 // must go to a real shell (redirection, expansions, escapes). Bail out
580 // before tokenizing so we neither waste work nor parse shell syntax we
581 // deliberately delegate.
582 if cmd_name.is_empty() {
583 return None;
584 }
585
586 // Parse args from command string, respecting quotes and performing
587 // POSIX quote removal so `echo label:'help wanted'` reaches the built-in
588 // as the single argument `label:help wanted` (issue #48).
589 let words = shell_parser::split_command_words(&self.command);
590 let args: Vec<String> = words.into_iter().skip(1).collect();
591
592 let ctx = CommandContext {
593 args,
594 stdin: match &self.options.stdin {
595 StdinOption::Content(s) => Some(s.clone()),
596 _ => None,
597 },
598 cwd: self.options.cwd.clone(),
599 env: self.options.env.clone(),
600 output_tx: self.output_tx.clone(),
601 is_cancelled: None,
602 };
603
604 match cmd_name {
605 "echo" => Some(commands::echo(ctx).await),
606 "pwd" => Some(commands::pwd(ctx).await),
607 "cd" => Some(commands::cd::resolve_cd(ctx).await.0),
608 "true" => Some(commands::r#true(ctx).await),
609 "false" => Some(commands::r#false(ctx).await),
610 "sleep" => Some(commands::sleep(ctx).await),
611 "cat" => Some(commands::cat(ctx).await),
612 "ls" => Some(commands::ls(ctx).await),
613 "mkdir" => Some(commands::mkdir(ctx).await),
614 "rm" => Some(commands::rm(ctx).await),
615 "touch" => Some(commands::touch(ctx).await),
616 "cp" => Some(commands::cp(ctx).await),
617 "mv" => Some(commands::mv(ctx).await),
618 "basename" => Some(commands::basename(ctx).await),
619 "dirname" => Some(commands::dirname(ctx).await),
620 "env" => Some(commands::env(ctx).await),
621 "exit" => Some(commands::exit(ctx).await),
622 "which" => Some(commands::which(ctx).await),
623 "yes" => Some(commands::yes(ctx).await),
624 "seq" => Some(commands::seq(ctx).await),
625 "tee" => Some(commands::tee(ctx).await),
626 "test" => Some(commands::test(ctx).await),
627 "head" => Some(commands::head(ctx).await),
628 "tail" => Some(commands::tail(ctx).await),
629 "sort" => Some(commands::sort(ctx).await),
630 "uniq" => Some(commands::uniq(ctx).await),
631
632 _ => None,
633 }
634 }
635
636 /// Stop the process using the configured kill signal
637 /// ([`RunOptions::kill_signal`], default `SIGTERM`).
638 ///
639 /// Mirrors the JavaScript `kill()` with no argument.
640 pub fn kill(&mut self) -> Result<()> {
641 let signal = self.options.kill_signal.clone();
642 self.kill_with(&signal)
643 }
644
645 /// Stop the process using an explicit signal, overriding
646 /// [`RunOptions::kill_signal`] for this call.
647 ///
648 /// Mirrors the JavaScript `kill(signal)`. The signal is delivered to the
649 /// child and its process group, so grandchildren behind a `sh -c` wrapper
650 /// are stopped too - except for a child sharing the caller's terminal,
651 /// which stays in the caller's group so CTRL+C keeps reaching it. The child
652 /// then has [`RunOptions::kill_grace_ms`] to run its own handler before
653 /// `SIGKILL` follows, so a process that ignores the signal still
654 /// terminates.
655 ///
656 /// ```no_run
657 /// use command_stream::{ProcessRunner, RunOptions};
658 ///
659 /// # #[tokio::main]
660 /// # async fn main() -> command_stream::Result<()> {
661 /// let mut runner = ProcessRunner::new("sleep 30", RunOptions::default());
662 /// runner.start().await?;
663 /// runner.kill_with("SIGINT")?; // the CTRL+C signal
664 /// # Ok(())
665 /// # }
666 /// ```
667 pub fn kill_with(&mut self, signal: &str) -> Result<()> {
668 self.cancelled = true;
669 utils::trace_lazy("ProcessRunner", || format!("kill | signal={signal}"));
670
671 let Some(child) = self.child.as_mut() else {
672 return Ok(());
673 };
674
675 // Windows stops the tree before the parent exits, while its descendants
676 // can still be found. Direct-child termination is the fallback.
677 // The `#[cfg(unix)]` block below is stripped on Windows, which leaves
678 // this one as the function's tail expression - hence no `return`.
679 #[cfg(not(unix))]
680 {
681 if let Some(pid) = child.id() {
682 signal::send_signal_to_process(pid, signal, signal::Delivery::ProcessAndGroup);
683 }
684 child.start_kill()?;
685 Ok(())
686 }
687
688 // Without a pid the process never spawned (or was already reaped);
689 // fall back to the forceful stop so `kill()` still terminates it.
690 #[cfg(unix)]
691 {
692 let Some(pid) = child.id() else {
693 child.start_kill()?;
694 return Ok(());
695 };
696
697 // `SIGKILL` cannot be handled, so there is nothing to wait for.
698 //
699 // A zero grace period means the child is given no opportunity to
700 // handle the signal either, so the requested signal is not
701 // delivered at all. Anything done between it and `SIGKILL` - even a
702 // single syscall - is a window the child can be scheduled in, which
703 // made "no grace" a race the child occasionally won rather than a
704 // guarantee. The reported exit code still comes from the signal
705 // that was requested.
706 let grace = self.options.kill_grace_ms;
707 let delivery = if self.own_process_group {
708 signal::Delivery::ProcessAndGroup
709 } else {
710 signal::Delivery::ProcessOnly
711 };
712 if grace == 0 || signal == "SIGKILL" {
713 signal::send_signal_to_process(pid, "SIGKILL", delivery);
714 let _ = child.start_kill();
715 return Ok(());
716 }
717
718 signal::send_signal_to_process(pid, signal, delivery);
719
720 // Otherwise escalate in the background so the child keeps its grace
721 // period without blocking the caller, which may not be inside an
722 // await point.
723 tokio::spawn(async move {
724 tokio::time::sleep(std::time::Duration::from_millis(grace)).await;
725 // Best effort: if the child already exited on the first signal
726 // this delivery simply fails, and the pid has not been reused
727 // because the `Child` handle above has not reaped it yet. That
728 // unreaped leader is also what keeps the group id alive, so the
729 // group delivery still reaches a grandchild that outlived it.
730 signal::send_signal_to_process(pid, "SIGKILL", delivery);
731 });
732
733 Ok(())
734 }
735 }
736
737 /// Check if the process is finished
738 pub fn is_finished(&self) -> bool {
739 self.finished
740 }
741
742 /// Get the result if available
743 pub fn result(&self) -> Option<&CommandResult> {
744 self.result.as_ref()
745 }
746
747 /// Process id of the command, or `None` when there is no operating system
748 /// process to identify.
749 ///
750 /// It is `None` before the command starts, and stays `None` for built-in
751 /// (virtual) commands such as `echo` or `sleep`, which run inside this
752 /// process and never spawn a child. Once a real command has been spawned
753 /// the value is stable: it remains readable after the command finishes,
754 /// unlike the child handle, which [`run`](Self::run) consumes.
755 ///
756 /// Mirrors the JavaScript `runner.pid` property.
757 ///
758 /// ```no_run
759 /// use command_stream::{ProcessRunner, RunOptions};
760 ///
761 /// # #[tokio::main]
762 /// # async fn main() -> command_stream::Result<()> {
763 /// let mut runner = ProcessRunner::new("/bin/sleep 1", RunOptions::default());
764 /// runner.start().await?;
765 /// println!("running as pid {:?}", runner.pid());
766 /// runner.run().await?;
767 /// println!("still readable: {:?}", runner.pid());
768 /// # Ok(())
769 /// # }
770 /// ```
771 pub fn pid(&self) -> Option<u32> {
772 self.pid
773 }
774
775 /// Get the command string
776 pub fn command(&self) -> &str {
777 &self.command
778 }
779
780 /// Get the options
781 pub fn options(&self) -> &RunOptions {
782 &self.options
783 }
784}
785
786/// Execute a command and return the result
787///
788/// This is the main entry point for simple command execution.
789/// Named `run` instead of `$` since `$` is not a valid Rust identifier.
790pub async fn run(command: impl Into<String>) -> Result<CommandResult> {
791 let mut runner = ProcessRunner::new(command, RunOptions::default());
792 runner.run().await
793}
794
795/// Alias for `run` function - for JavaScript-like API feel
796/// Since `$` is not valid in Rust, this provides a similar short name
797pub use run as execute;
798
799/// Execute a command with custom options
800pub async fn exec(command: impl Into<String>, options: RunOptions) -> Result<CommandResult> {
801 let mut runner = ProcessRunner::new(command, options);
802 runner.run().await
803}
804
805/// Create a new process runner without starting it
806pub fn create(command: impl Into<String>, options: RunOptions) -> ProcessRunner {
807 ProcessRunner::new(command, options)
808}
809
810/// Execute a command synchronously (blocking)
811pub fn run_sync(command: impl Into<String>) -> Result<CommandResult> {
812 let rt = tokio::runtime::Runtime::new()?;
813 rt.block_on(run(command))
814}
815
816// Tests are located in tests/ directory for better organization