command_stream/stream.rs
1//! Streaming and async iteration support
2//!
3//! This module provides async streaming capabilities similar to JavaScript's
4//! async iterators and stream handling in `$.stream-utils.mjs`.
5//!
6//! It mirrors the JavaScript implementation's behavior for issue #155:
7//!
8//! 1. The stream yields an explicit `OutputChunk::Exit(code)` when the
9//! process exits, so consumers can observe the exit code from inside the
10//! loop.
11//! 2. The stream does not hang forever when the process has exited but a
12//! grandchild keeps the stdio pipes open (the readers are drained with a
13//! grace period and then aborted).
14//! 3. The process can be stopped from inside the loop via
15//! [`OutputStream::kill`] / [`OutputStream::kill_with`], and abandoning the
16//! stream (e.g. `break`) also stops the process.
17//! 4. The stop signal is configurable via
18//! [`StreamingRunner::kill_signal`] (default `SIGTERM`), just like the
19//! JavaScript `killSignal` option.
20//!
21//! ## Usage
22//!
23//! ```rust,no_run
24//! use command_stream::{StreamingRunner, OutputChunk};
25//!
26//! #[tokio::main]
27//! async fn main() -> Result<(), Box<dyn std::error::Error>> {
28//! let runner = StreamingRunner::new("yes hello");
29//!
30//! // Stream output as it arrives
31//! let mut stream = runner.stream();
32//! let mut count = 0;
33//! while let Some(chunk) = stream.next().await {
34//! match chunk {
35//! OutputChunk::Stdout(data) => {
36//! print!("{}", String::from_utf8_lossy(&data));
37//! count += 1;
38//! if count >= 5 {
39//! // Stop the process from inside the loop.
40//! stream.kill();
41//! }
42//! }
43//! OutputChunk::Stderr(data) => {
44//! eprint!("{}", String::from_utf8_lossy(&data));
45//! }
46//! OutputChunk::Exit(code) => {
47//! println!("Process exited with code: {}", code);
48//! break;
49//! }
50//! }
51//! }
52//!
53//! Ok(())
54//! }
55//! ```
56
57use std::collections::HashMap;
58use std::ffi::OsString;
59use std::path::PathBuf;
60use std::process::Stdio;
61use std::time::Duration;
62use tokio::io::BufReader;
63use tokio::process::Command;
64use tokio::sync::{mpsc, watch};
65use tokio::task::JoinHandle;
66
67use crate::signal::{
68 send_signal_to_process, signal_exit_code, Delivery, DEFAULT_KILL_GRACE_MS, DEFAULT_KILL_SIGNAL,
69};
70use crate::trace::trace_lazy;
71use crate::{CommandResult, Result};
72
73/// Default grace period (in milliseconds) to keep draining the stdio pipes
74/// after the process has exited before aborting any lingering readers. Mirrors
75/// the JavaScript `exitPumpGrace` default.
76const DEFAULT_EXIT_PUMP_GRACE_MS: u64 = 100;
77
78/// A chunk of output from a streaming process
79#[derive(Debug, Clone)]
80pub enum OutputChunk {
81 /// Stdout data
82 Stdout(Vec<u8>),
83 /// Stderr data
84 Stderr(Vec<u8>),
85 /// Process exit code
86 Exit(i32),
87}
88
89/// A streaming process runner that allows async iteration over output
90pub struct StreamingRunner {
91 command: StreamingCommand,
92 cwd: Option<PathBuf>,
93 env: Option<HashMap<String, String>>,
94 prefer_local: crate::PreferLocal,
95 stdin_content: Option<String>,
96 kill_signal: String,
97 kill_grace_ms: u64,
98 exit_pump_grace_ms: u64,
99}
100
101#[derive(Clone)]
102enum StreamingCommand {
103 Shell(String),
104 Argv {
105 program: OsString,
106 args: Vec<OsString>,
107 },
108}
109
110impl StreamingRunner {
111 /// Create a streaming runner for a command string interpreted by the
112 /// platform shell.
113 pub fn new(command: impl Into<String>) -> Self {
114 Self::with_command(StreamingCommand::Shell(command.into()))
115 }
116
117 /// Create a streaming runner for an executable and exact argument vector.
118 ///
119 /// Unlike [`StreamingRunner::new`], this constructor bypasses the platform
120 /// shell. Argument boundaries are therefore preserved on every platform,
121 /// including Windows, without requiring shell-specific quoting.
122 pub fn from_argv<P, I, S>(program: P, args: I) -> Self
123 where
124 P: Into<OsString>,
125 I: IntoIterator<Item = S>,
126 S: Into<OsString>,
127 {
128 Self::with_command(StreamingCommand::Argv {
129 program: program.into(),
130 args: args.into_iter().map(Into::into).collect(),
131 })
132 }
133
134 fn with_command(command: StreamingCommand) -> Self {
135 StreamingRunner {
136 command,
137 cwd: None,
138 env: None,
139 prefer_local: crate::PreferLocal::Off,
140 stdin_content: None,
141 kill_signal: DEFAULT_KILL_SIGNAL.to_string(),
142 kill_grace_ms: DEFAULT_KILL_GRACE_MS,
143 exit_pump_grace_ms: DEFAULT_EXIT_PUMP_GRACE_MS,
144 }
145 }
146
147 /// Set the working directory
148 pub fn cwd(mut self, path: impl Into<PathBuf>) -> Self {
149 self.cwd = Some(path.into());
150 self
151 }
152
153 /// Set environment variables
154 pub fn env(mut self, env: HashMap<String, String>) -> Self {
155 self.env = Some(env);
156 self
157 }
158
159 /// Prefer executables from the command's working directory or explicit directories.
160 pub fn prefer_local(mut self, preference: crate::PreferLocal) -> Self {
161 self.prefer_local = preference;
162 self
163 }
164
165 /// Set stdin content
166 pub fn stdin(mut self, content: impl Into<String>) -> Self {
167 self.stdin_content = Some(content.into());
168 self
169 }
170
171 /// Configure the signal used to stop the process when it is killed without
172 /// an explicit signal — i.e. [`OutputStream::kill`] or abandoning the
173 /// stream. Mirrors the JavaScript `killSignal` option (default `SIGTERM`).
174 ///
175 /// The reported exit code follows the conventional `128 + signal` mapping
176 /// (e.g. `SIGTERM` => 143, `SIGINT` => 130, `SIGKILL` => 137).
177 pub fn kill_signal(mut self, signal: impl Into<String>) -> Self {
178 self.kill_signal = signal.into();
179 self
180 }
181
182 /// Configure how long (in milliseconds) the child is given to handle the
183 /// kill signal before `SIGKILL` is sent. Mirrors the JavaScript `killGrace`
184 /// option (default 100ms).
185 ///
186 /// This is the window in which a child running its own `SIGTERM` handler
187 /// can shut down on its own terms. Set it to `0` to escalate immediately.
188 pub fn kill_grace_ms(mut self, ms: u64) -> Self {
189 self.kill_grace_ms = ms;
190 self
191 }
192
193 /// Configure the grace period (in milliseconds) to keep draining the stdio
194 /// pipes after the process exits before aborting lingering readers. Mirrors
195 /// the JavaScript `exitPumpGrace` option (default 100ms).
196 pub fn exit_pump_grace_ms(mut self, ms: u64) -> Self {
197 self.exit_pump_grace_ms = ms;
198 self
199 }
200
201 fn spawn(mut self) -> (OutputStream, JoinHandle<Result<()>>) {
202 let (tx, rx) = mpsc::channel(1024);
203 // Unbounded so a synchronous Drop can request a kill without awaiting.
204 let (kill_tx, kill_rx) = mpsc::unbounded_channel::<String>();
205 // The child is spawned inside the task below, so its id is not known
206 // when this returns. The task publishes it here as soon as the spawn
207 // succeeds; `OutputStream::pid` reads the latest value (issue #18).
208 let (pid_tx, pid_rx) = watch::channel(None);
209
210 // Spawn the process handling task
211 let command = self.command.clone();
212 let cwd = self.cwd.take();
213 let mut env = self.env.take();
214 let local_cwd = cwd
215 .clone()
216 .or_else(|| std::env::current_dir().ok())
217 .unwrap_or_else(|| PathBuf::from("."));
218 if let Some((key, path)) =
219 crate::local_bin::preferred_path(env.as_ref(), &local_cwd, &self.prefer_local)
220 {
221 env.get_or_insert_with(|| std::env::vars().collect())
222 .insert(key, path);
223 }
224 let stdin_content = self.stdin_content.take();
225 let grace = GraceWindows {
226 exit_pump_ms: self.exit_pump_grace_ms,
227 kill_ms: self.kill_grace_ms,
228 };
229 let kill_signal = self.kill_signal.clone();
230
231 let task = tokio::spawn(async move {
232 let channels = StreamChannels {
233 output_tx: tx,
234 kill_rx,
235 pid_tx,
236 };
237 let result =
238 run_streaming_process(command, cwd, env, stdin_content, grace, channels).await;
239 if let Err(error) = &result {
240 trace_lazy("StreamingRunner", || format!("Error: {error}"));
241 }
242 result
243 });
244
245 (
246 OutputStream {
247 rx,
248 kill_tx,
249 kill_signal,
250 killed: false,
251 pid_rx,
252 },
253 task,
254 )
255 }
256
257 /// Start the process and return a stream of output chunks
258 pub fn stream(self) -> OutputStream {
259 self.spawn().0
260 }
261
262 /// Run to completion and collect all output
263 pub async fn collect(self) -> Result<CommandResult> {
264 let stdin_content = self.stdin_content.clone();
265 let mut stdout = Vec::new();
266 let mut stderr = Vec::new();
267 let mut exit_code = 0;
268
269 let (mut stream, task) = self.spawn();
270 while let Some(chunk) = stream.rx.recv().await {
271 match chunk {
272 OutputChunk::Stdout(data) => stdout.extend(data),
273 OutputChunk::Stderr(data) => stderr.extend(data),
274 OutputChunk::Exit(code) => exit_code = code,
275 }
276 }
277
278 task.await.map_err(|error| {
279 std::io::Error::other(format!("streaming process task failed: {error}"))
280 })??;
281
282 let mut result = CommandResult::new(
283 String::from_utf8_lossy(&stdout).to_string(),
284 String::from_utf8_lossy(&stderr).to_string(),
285 exit_code,
286 );
287 if let Some(content) = stdin_content {
288 result.stdin = crate::result_streams::CapturedInput::new(content.into_bytes());
289 }
290 Ok(result)
291 }
292
293 /// Run an exact-argument command to completion from synchronous code.
294 ///
295 /// Build the command with [`Self::from_argv`] and configure it with the
296 /// same `cwd`, `env`, and `stdin` methods used by [`Self::collect`]. Call
297 /// this outside a Tokio runtime; async callers should use `collect().await`.
298 pub fn collect_blocking(self) -> Result<CommandResult> {
299 if tokio::runtime::Handle::try_current().is_ok() {
300 return Err(std::io::Error::other(
301 "collect_blocking cannot run inside a Tokio runtime; use collect().await",
302 )
303 .into());
304 }
305 let runtime = tokio::runtime::Runtime::new()?;
306 runtime.block_on(self.collect())
307 }
308}
309
310/// Stream of output chunks from a process
311pub struct OutputStream {
312 rx: mpsc::Receiver<OutputChunk>,
313 kill_tx: mpsc::UnboundedSender<String>,
314 kill_signal: String,
315 killed: bool,
316 pid_rx: watch::Receiver<Option<u32>>,
317}
318
319impl OutputStream {
320 /// Receive the next chunk
321 pub async fn next(&mut self) -> Option<OutputChunk> {
322 self.rx.recv().await
323 }
324
325 /// Process id of the streamed command, as currently known.
326 ///
327 /// The child is spawned by a background task, so this is `None` for the
328 /// short window between [`StreamingRunner::stream`] returning and the spawn
329 /// completing, and stays `None` if the spawn failed. From the first
330 /// delivered chunk onwards it is set, and it remains readable after the
331 /// process has exited. Use [`wait_for_pid`](Self::wait_for_pid) to avoid
332 /// the startup window.
333 pub fn pid(&self) -> Option<u32> {
334 *self.pid_rx.borrow()
335 }
336
337 /// Process id of the streamed command, waiting for the spawn to complete.
338 ///
339 /// Resolves as soon as the child exists, and returns `None` if the process
340 /// could never be spawned. This is the streaming counterpart of awaiting a
341 /// stream before reading `runner.pid` in JavaScript.
342 pub async fn wait_for_pid(&mut self) -> Option<u32> {
343 // `wait_for` checks the current value first, so an already-published id
344 // returns without waiting. An error means the sending task is gone,
345 // which only happens when the spawn failed.
346 match self.pid_rx.wait_for(|pid| pid.is_some()).await {
347 Ok(pid) => *pid,
348 Err(_) => None,
349 }
350 }
351
352 /// Stop the process using the configured kill signal (default `SIGTERM`).
353 ///
354 /// This can be called from inside the consumption loop to stop a
355 /// long-running or endless process; a terminating `OutputChunk::Exit` is
356 /// still delivered afterwards.
357 pub fn kill(&mut self) {
358 let signal = self.kill_signal.clone();
359 self.kill_with(&signal);
360 }
361
362 /// Stop the process using an explicit signal, overriding the configured
363 /// kill signal for this call.
364 pub fn kill_with(&mut self, signal: &str) {
365 if self.killed {
366 return;
367 }
368 self.killed = true;
369 trace_lazy("OutputStream", || format!("kill | signal={}", signal));
370 // Best effort: the task may have already finished, in which case the
371 // receiver is gone and the send fails harmlessly.
372 let _ = self.kill_tx.send(signal.to_string());
373 }
374
375 /// Collect all remaining output into vectors
376 pub async fn collect(mut self) -> (Vec<u8>, Vec<u8>, i32) {
377 let mut stdout = Vec::new();
378 let mut stderr = Vec::new();
379 let mut exit_code = 0;
380
381 while let Some(chunk) = self.rx.recv().await {
382 match chunk {
383 OutputChunk::Stdout(data) => stdout.extend(data),
384 OutputChunk::Stderr(data) => stderr.extend(data),
385 OutputChunk::Exit(code) => exit_code = code,
386 }
387 }
388
389 (stdout, stderr, exit_code)
390 }
391
392 /// Collect stdout only, discarding stderr
393 pub async fn collect_stdout(mut self) -> Vec<u8> {
394 let mut stdout = Vec::new();
395
396 while let Some(chunk) = self.rx.recv().await {
397 if let OutputChunk::Stdout(data) = chunk {
398 stdout.extend(data);
399 }
400 }
401
402 stdout
403 }
404}
405
406impl Drop for OutputStream {
407 fn drop(&mut self) {
408 // Abandoning the stream (e.g. `break`-ing out of the loop) must stop the
409 // process, matching the JavaScript iterator's `finally` cleanup. If the
410 // process already finished this is a harmless no-op.
411 if !self.killed {
412 let _ = self.kill_tx.send(self.kill_signal.clone());
413 }
414 }
415}
416
417/// The channels `run_streaming_process` communicates over: output chunks out,
418/// kill requests in, and the child's id published once the spawn succeeds.
419struct StreamChannels {
420 /// Carries the output chunks, and finally the `Exit` chunk, to the consumer.
421 output_tx: mpsc::Sender<OutputChunk>,
422 /// Carries kill requests, by signal name, in from the consumer.
423 kill_rx: mpsc::UnboundedReceiver<String>,
424 /// Publishes the child's id, which is only known inside the spawning task.
425 pid_tx: watch::Sender<Option<u32>>,
426}
427
428/// How long the runner waits, in milliseconds, at the two points where it gives
429/// something a chance to finish on its own before forcing the issue.
430#[derive(Debug, Clone, Copy)]
431struct GraceWindows {
432 /// Time allowed for the readers to drain buffered output after the child
433 /// exits, before the `Exit` chunk is emitted.
434 exit_pump_ms: u64,
435 /// Time allowed for the child to handle the delivered signal, before the
436 /// escalation to `SIGKILL`.
437 kill_ms: u64,
438}
439
440/// Run a streaming process and send output to the channel
441async fn run_streaming_process(
442 command: StreamingCommand,
443 cwd: Option<PathBuf>,
444 env: Option<HashMap<String, String>>,
445 stdin_content: Option<String>,
446 grace: GraceWindows,
447 channels: StreamChannels,
448) -> Result<()> {
449 let StreamChannels {
450 output_tx: tx,
451 mut kill_rx,
452 pid_tx,
453 } = channels;
454 trace_lazy("StreamingRunner", || match &command {
455 StreamingCommand::Shell(command) => format!("Starting: {command}"),
456 StreamingCommand::Argv { program, args } => {
457 format!("Starting argv command: {program:?} {args:?}")
458 }
459 });
460
461 let mut cmd = match command {
462 StreamingCommand::Shell(command) => crate::utils::shell_command(&command, env.as_ref()),
463 StreamingCommand::Argv { program, args } => {
464 let mut cmd = Command::new(program);
465 cmd.args(args);
466 cmd
467 }
468 };
469
470 // Configure stdio
471 if stdin_content.is_some() {
472 cmd.stdin(Stdio::piped());
473 } else {
474 cmd.stdin(Stdio::null());
475 }
476 cmd.stdout(Stdio::piped());
477 cmd.stderr(Stdio::piped());
478
479 // Run the child in its own process group so we can signal the whole group
480 // (parent + grandchildren), matching the JavaScript implementation.
481 #[cfg(unix)]
482 cmd.process_group(0);
483
484 // Set working directory
485 if let Some(ref cwd) = cwd {
486 cmd.current_dir(cwd);
487 }
488
489 // Set environment
490 if let Some(ref env_vars) = env {
491 for (key, value) in env_vars {
492 cmd.env(key, value);
493 }
494 }
495
496 let mut child = cmd.spawn()?;
497 // Publish the id before any awaiting, so a consumer asking for it as soon
498 // as the first chunk arrives already sees it.
499 let _ = pid_tx.send(child.id());
500
501 // Write stdin if needed
502 if let Some(content) = stdin_content {
503 if let Some(mut stdin) = child.stdin.take() {
504 use tokio::io::AsyncWriteExt;
505 let _ = stdin.write_all(content.as_bytes()).await;
506 let _ = stdin.shutdown().await;
507 }
508 }
509
510 // Spawn stdout reader
511 let stdout = child.stdout.take();
512 let tx_stdout = tx.clone();
513 let stdout_handle = stdout.map(|stdout| {
514 tokio::spawn(async move {
515 let mut reader = BufReader::new(stdout);
516 let mut buf = vec![0u8; 8192];
517 loop {
518 use tokio::io::AsyncReadExt;
519 match reader.read(&mut buf).await {
520 Ok(0) => break,
521 Ok(n) => {
522 if tx_stdout
523 .send(OutputChunk::Stdout(buf[..n].to_vec()))
524 .await
525 .is_err()
526 {
527 break;
528 }
529 }
530 Err(_) => break,
531 }
532 }
533 })
534 });
535
536 // Spawn stderr reader
537 let stderr = child.stderr.take();
538 let tx_stderr = tx.clone();
539 let stderr_handle = stderr.map(|stderr| {
540 tokio::spawn(async move {
541 let mut reader = BufReader::new(stderr);
542 let mut buf = vec![0u8; 8192];
543 loop {
544 use tokio::io::AsyncReadExt;
545 match reader.read(&mut buf).await {
546 Ok(0) => break,
547 Ok(n) => {
548 if tx_stderr
549 .send(OutputChunk::Stderr(buf[..n].to_vec()))
550 .await
551 .is_err()
552 {
553 break;
554 }
555 }
556 Err(_) => break,
557 }
558 }
559 })
560 });
561
562 // Wait for the process to exit OR for a kill request — crucially we do NOT
563 // wait for the readers first. If a grandchild keeps the pipe open the
564 // readers would never finish, so waiting on them before the exit (as the
565 // old implementation did) would hang forever (issue #155).
566 let pid = child.id();
567 let code;
568 tokio::select! {
569 status = child.wait() => {
570 code = status_to_code(status?);
571 }
572 maybe_signal = kill_rx.recv() => {
573 // A kill was requested (explicit kill()/kill_with() or the stream
574 // being dropped). Stop the process group with the requested signal.
575 let signal = maybe_signal.unwrap_or_else(|| DEFAULT_KILL_SIGNAL.to_string());
576 trace_lazy("StreamingRunner", || format!("Kill requested | signal={}", signal));
577 // Give the child its grace period to run its own handler and exit
578 // on its own terms, then escalate to a forceful kill so a process
579 // that ignores the signal still terminates.
580 //
581 // A zero grace period means the child is given no opportunity to
582 // handle the signal, so the requested signal is not delivered at
583 // all. Anything done between it and the forceful kill - a syscall,
584 // or awaiting a zero-length timeout, which yields to the runtime -
585 // is a window the child can be scheduled in, which made "no grace"
586 // a race the child occasionally won rather than a guarantee.
587 let survived_grace = if grace.kill_ms == 0 {
588 true
589 } else {
590 if let Some(pid) = pid {
591 // The child is always spawned with `process_group(0)`
592 // above, so it leads the group named by its own pid.
593 send_signal_to_process(pid, &signal, Delivery::ProcessAndGroup);
594 }
595 tokio::time::timeout(Duration::from_millis(grace.kill_ms), child.wait())
596 .await
597 .is_err()
598 };
599 if survived_grace {
600 if let Some(pid) = pid {
601 send_signal_to_process(pid, "SIGKILL", Delivery::ProcessAndGroup);
602 }
603 let _ = child.start_kill();
604 let _ = child.wait().await;
605 }
606 // Report the conventional 128 + signal code for the requested
607 // signal, matching the JavaScript implementation.
608 code = signal_exit_code(&signal);
609 }
610 }
611
612 // The process has exited. Give the readers a short grace period to flush any
613 // buffered output, then abort any that are still blocked on an inherited
614 // open pipe so we don't hang.
615 let stdout_abort = stdout_handle.as_ref().map(|h| h.abort_handle());
616 let stderr_abort = stderr_handle.as_ref().map(|h| h.abort_handle());
617 let drain = async {
618 if let Some(handle) = stdout_handle {
619 let _ = handle.await;
620 }
621 if let Some(handle) = stderr_handle {
622 let _ = handle.await;
623 }
624 };
625 if tokio::time::timeout(Duration::from_millis(grace.exit_pump_ms), drain)
626 .await
627 .is_err()
628 {
629 // A reader is still blocked on an inherited open pipe — abort it so the
630 // exit chunk is delivered without waiting for the grandchild.
631 if let Some(abort) = stdout_abort {
632 abort.abort();
633 }
634 if let Some(abort) = stderr_abort {
635 abort.abort();
636 }
637 }
638
639 // Send exit code (always — even if a reader was aborted).
640 let _ = tx.send(OutputChunk::Exit(code)).await;
641
642 trace_lazy("StreamingRunner", || format!("Exited with code: {}", code));
643
644 Ok(())
645}
646
647/// Convert an exit status into a numeric exit code, using the conventional
648/// `128 + signal` mapping when the process was terminated by a signal.
649fn status_to_code(status: std::process::ExitStatus) -> i32 {
650 if let Some(code) = status.code() {
651 return code;
652 }
653 #[cfg(unix)]
654 {
655 use std::os::unix::process::ExitStatusExt;
656 if let Some(sig) = status.signal() {
657 return 128 + sig;
658 }
659 }
660 -1
661}
662
663/// Async iterator trait for output streams
664#[async_trait::async_trait]
665pub trait AsyncIterator {
666 type Item;
667
668 /// Get the next item from the iterator
669 async fn next(&mut self) -> Option<Self::Item>;
670}
671
672#[async_trait::async_trait]
673impl AsyncIterator for OutputStream {
674 type Item = OutputChunk;
675
676 async fn next(&mut self) -> Option<Self::Item> {
677 self.rx.recv().await
678 }
679}
680
681/// Extension trait to convert ProcessRunner into a stream
682pub trait IntoStream {
683 /// Convert into an output stream
684 fn into_stream(self) -> OutputStream;
685}
686
687impl IntoStream for crate::ProcessRunner {
688 fn into_stream(self) -> OutputStream {
689 let streaming = StreamingRunner::new(self.command().to_string());
690 streaming.stream()
691 }
692}