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