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