command-stream 1.2.0

Modern shell command execution library with streaming, async iteration, and event support
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
//! The Bun Shell interpreter: Bun's `interpreter.rs` and its state machines
//! (`src/runtime/shell/states/`: Script, Stmt, Binary, Pipeline, Subshell,
//! If, CondExpr, Assigns, Cmd; MIT, Copyright (c) Oven-sh / Jarred Sumner),
//! ported by way of `js/src/bun-shell/interpreter.mjs`. Each state becomes an
//! async method that resolves to the node's exit code.
//!
//! Errors that Bun throws into JS (a failed write of an error message, an
//! invalid JS redirect target) reject the whole run: they are the `Err` of
//! [`ExecResult`]. Everything else is reported on stderr with an exit code.
//!
//! Builtins borrow the [`ShellExecEnv`] mutably, so the futures are not
//! `'static`: pipeline items are driven concurrently on the caller's task
//! (see [`join_all`]) instead of being spawned.

use std::future::Future;
use std::pin::Pin;
use std::task::Poll;

use super::env::{EnvKind, EnvMap, ShellExecEnv};
use super::expansion::{expand_atom, ExpandError, ExpandOpts, Expanded};
use super::io::{Channel, InKind, OutKind, Reader, ShellIO, ShellSysError, Writer};
use super::parser::{Atom, Binary, BinaryOp, CondExpr, CondExprOp, Expr, If, Script, Stmt};
use super::{ShellError, ShellValue};

mod cmd;

/// A boxed, sendable future (breaks the async recursion).
pub(crate) type BoxFut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;

/// An exit code, or the error that rejects the run.
pub(crate) type ExecResult = Result<i32, ShellError>;

/// How to set up an [`Interpreter`] (the JS constructor options).
pub(crate) struct InterpreterOptions {
    /// JS values referenced by the script (`\x08__bun_N\x08`).
    pub(crate) jsobjs: Vec<ShellValue>,
    /// The export environment.
    pub(crate) env: EnvMap,
    /// Initial directory (default: the process cwd).
    pub(crate) cwd: Option<String>,
    /// Buffer stdout/stderr instead of echoing them.
    pub(crate) quiet: bool,
    /// Positional parameters (`$0..$9`).
    pub(crate) argv: Vec<String>,
}

/// The result of a finished script.
#[derive(Debug)]
pub(crate) struct RunOutput {
    pub(crate) exit_code: i32,
    pub(crate) stdout: Vec<u8>,
    pub(crate) stderr: Vec<u8>,
}

/// The shared, read-only part of a running script.
pub(crate) struct Interpreter {
    pub(crate) jsobjs: Vec<ShellValue>,
    pub(crate) argv: Vec<String>,
    root_io: ShellIO,
}

/// Poll every future until all are done (`Promise.all` without the early
/// rejection); results keep the input order.
pub(crate) async fn join_all<'a, T: Send + 'a>(futs: Vec<BoxFut<'a, T>>) -> Vec<T> {
    let mut futs: Vec<Option<BoxFut<'a, T>>> = futs.into_iter().map(Some).collect();
    let mut results: Vec<Option<T>> = futs.iter().map(|_| None).collect();
    std::future::poll_fn(|cx| {
        let mut pending = false;
        for (slot, result) in futs.iter_mut().zip(results.iter_mut()) {
            if let Some(f) = slot {
                match f.as_mut().poll(cx) {
                    Poll::Ready(v) => {
                        *result = Some(v);
                        *slot = None;
                    }
                    Poll::Pending => pending = true,
                }
            }
        }
        if pending {
            Poll::Pending
        } else {
            Poll::Ready(())
        }
    })
    .await;
    results.into_iter().flatten().collect()
}

impl Interpreter {
    /// Create the interpreter and its root environment. Fails like the JS
    /// constructor when the cwd cannot be entered.
    pub(crate) fn new(opts: InterpreterOptions) -> Result<(Self, ShellExecEnv), ShellError> {
        let process_cwd = std::env::current_dir()
            .map(|p| p.to_string_lossy().into_owned())
            .unwrap_or_default();
        let mut root = ShellExecEnv::new(opts.env, process_cwd);
        if let Some(cwd) = &opts.cwd {
            root.change_cwd(cwd, true)
                .map_err(|e| ShellError::system(e.message()))?;
        }
        let out = |writer: fn() -> Writer, captured| {
            if opts.quiet {
                OutKind::Pipe
            } else {
                OutKind::fd(writer(), Some(captured))
            }
        };
        let root_io = ShellIO {
            stdin: InKind::Fd(Reader::stdin()),
            stdout: out(Writer::stdout, root.buffered_stdout.clone()),
            stderr: out(Writer::stderr, root.buffered_stderr.clone()),
        };
        Ok((
            Self {
                jsobjs: opts.jsobjs,
                argv: opts.argv,
                root_io,
            },
            root,
        ))
    }

    /// Run a parsed script in the root environment.
    pub(crate) async fn run(
        &self,
        script: &Script,
        root: &mut ShellExecEnv,
    ) -> Result<RunOutput, ShellError> {
        let io = self.root_io.clone();
        let exit_code = self.script(script, root, &io).await?;
        Ok(RunOutput {
            exit_code,
            stdout: root.buffered_stdout.to_vec(),
            stderr: root.buffered_stderr.to_vec(),
        })
    }

    // --- control flow ------------------------------------------------------

    pub(crate) async fn script(
        &self,
        node: &Script,
        shell: &mut ShellExecEnv,
        io: &ShellIO,
    ) -> ExecResult {
        self.stmts(&node.stmts, shell, io).await
    }

    async fn stmts(&self, stmts: &[Stmt], shell: &mut ShellExecEnv, io: &ShellIO) -> ExecResult {
        let mut exit_code = 0;
        for stmt in stmts {
            exit_code = 0;
            for expr in &stmt.exprs {
                exit_code = self.expr(expr, shell, io).await?;
            }
        }
        Ok(exit_code)
    }

    fn expr<'a>(
        &'a self,
        node: &'a Expr,
        shell: &'a mut ShellExecEnv,
        io: &'a ShellIO,
    ) -> BoxFut<'a, ExecResult> {
        Box::pin(async move {
            match node {
                Expr::Assign(assigns) => self.assigns(assigns, shell, false).await,
                Expr::Binary(b) => self.binary(b, shell, io).await,
                Expr::Pipeline(items) => self.pipeline(items, shell, io).await,
                Expr::Cmd(c) => self.cmd(c, shell, io).await,
                Expr::Subshell(s) => {
                    let mut env = shell.dupe_for_subshell(io, EnvKind::Subshell);
                    self.script(&s.script, &mut env, io).await
                }
                Expr::If(node) => self.if_clause(node, shell, io).await,
                Expr::CondExpr(node) => self.cond_expr(node, shell, io).await,
            }
        })
    }

    async fn binary(&self, node: &Binary, shell: &mut ShellExecEnv, io: &ShellIO) -> ExecResult {
        let left = self.expr(&node.left, shell, io).await?;
        if (node.op == BinaryOp::And && left != 0) || (node.op == BinaryOp::Or && left == 0) {
            return Ok(left);
        }
        self.expr(&node.right, shell, io).await
    }

    async fn pipeline(&self, items: &[Expr], shell: &ShellExecEnv, io: &ShellIO) -> ExecResult {
        let items: Vec<&Expr> = items
            .iter()
            .filter(|item| !matches!(item, Expr::Assign(_)))
            .collect();
        if items.is_empty() {
            return Ok(0);
        }
        let n = items.len();
        let channels: Vec<Channel> = (1..n).map(|_| Channel::new()).collect();
        let mut runs: Vec<BoxFut<'_, ExecResult>> = Vec::with_capacity(n);
        for (i, item) in items.into_iter().enumerate() {
            let stdin = if i == 0 {
                io.stdin.clone()
            } else {
                InKind::Fd(Reader::channel(channels[i - 1].clone()))
            };
            let stdout = if i == n - 1 {
                io.stdout.clone()
            } else {
                OutKind::fd(Writer::channel(channels[i].clone()), None)
            };
            let item_io = ShellIO {
                stdin,
                stdout,
                stderr: io.stderr.clone(),
            };
            let mut env = shell.dupe_for_subshell(&item_io, EnvKind::Pipeline);
            runs.push(Box::pin(async move {
                let result = self.expr(item, &mut env, &item_io).await;
                if i > 0 {
                    if let InKind::Fd(reader) = &item_io.stdin {
                        reader.close();
                    }
                }
                if i < n - 1 {
                    if let OutKind::Fd { writer, .. } = &item_io.stdout {
                        writer.close().await;
                    }
                }
                result
            }));
        }
        drop(channels);
        let mut last = 0;
        for code in join_all(runs).await {
            last = code?;
        }
        Ok(if n >= 2 { last } else { 0 })
    }

    async fn if_clause(&self, node: &If, shell: &mut ShellExecEnv, io: &ShellIO) -> ExecResult {
        if self.stmts(&node.cond, shell, io).await? == 0 {
            return self.stmts(&node.then, shell, io).await;
        }
        let parts = &node.else_parts;
        match parts.len() {
            0 => return Ok(0),
            1 => return self.stmts(&parts[0], shell, io).await,
            _ => {}
        }
        let mut i = 0;
        while i + 1 < parts.len() {
            if self.stmts(&parts[i], shell, io).await? == 0 {
                return self.stmts(&parts[i + 1], shell, io).await;
            }
            i += 2;
        }
        if i < parts.len() {
            self.stmts(&parts[i], shell, io).await
        } else {
            Ok(0)
        }
    }

    async fn cond_expr(&self, node: &CondExpr, shell: &ShellExecEnv, io: &ShellIO) -> ExecResult {
        let mut args = Vec::with_capacity(node.args.len());
        for atom in &node.args {
            match self.expand(atom, shell, ExpandOpts::default()).await {
                Ok(out) => args.push(out.buf),
                Err(e) => {
                    let msg = format!("{}\n", e.display()?);
                    return Ok(self.write_failing_error_no_throw(io, shell, &msg).await);
                }
            }
        }
        let first = args.first().map(String::as_str).unwrap_or("");
        Ok(match node.op {
            CondExprOp::IsFile | CondExprOp::IsDirectory | CondExprOp::IsCharDevice => {
                if first.is_empty() {
                    return Ok(1);
                }
                let Ok(st) = std::fs::metadata(shell.resolve(first)) else {
                    return Ok(1);
                };
                let ok = match node.op {
                    CondExprOp::IsFile => st.is_file(),
                    CondExprOp::IsDirectory => st.is_dir(),
                    _ => is_char_device(&st),
                };
                i32::from(!ok)
            }
            CondExprOp::IsEmpty => i32::from(!first.is_empty()),
            CondExprOp::IsNonEmpty => i32::from(first.is_empty()),
            CondExprOp::Eq => {
                i32::from(!(args.is_empty() || (args.len() >= 2 && args[0] == args[1])))
            }
            CondExprOp::NotEq => i32::from(!(args.len() >= 2 && args[0] != args[1])),
        })
    }

    // --- expansion ----------------------------------------------------------

    pub(crate) async fn expand(
        &self,
        atom: &Atom,
        shell: &ShellExecEnv,
        opts: ExpandOpts,
    ) -> Result<Expanded, ExpandError> {
        expand_atom(self, shell, atom, opts).await
    }

    /// `$(...)`: run in a child env whose stdout is buffered; resolves to the
    /// exit code and the (lossily decoded) output.
    pub(crate) fn cmd_subst<'a>(
        &'a self,
        script: &'a Script,
        shell: &'a ShellExecEnv,
    ) -> BoxFut<'a, Result<(i32, String), ShellError>> {
        Box::pin(async move {
            let io = ShellIO {
                stdin: self.root_io.stdin.clone(),
                stdout: OutKind::Pipe,
                stderr: self.root_io.stderr.clone(),
            };
            let mut env = shell.dupe_for_subshell(&io, EnvKind::CmdSubst);
            let exit_code = self.script(script, &mut env, &io).await?;
            let stdout = String::from_utf8_lossy(&env.buffered_stdout.to_vec()).into_owned();
            Ok((exit_code, stdout))
        })
    }

    /// Assignments (Bun's `Assigns` state): 0, or 1 on an expansion error
    /// (Bun reports nothing here).
    pub(crate) async fn assigns(
        &self,
        assigns: &[super::parser::Assign],
        shell: &mut ShellExecEnv,
        cmd_local: bool,
    ) -> ExecResult {
        for assign in assigns {
            let opts = ExpandOpts {
                is_assign: true,
                assign_ctx: true,
            };
            let out = match self.expand(&assign.value, shell, opts).await {
                Ok(out) => out,
                Err(e) => {
                    e.display()?;
                    return Ok(1);
                }
            };
            shell.assign_var(&assign.label, out.words().join(" "), cmd_local);
        }
        Ok(0)
    }

    // --- errors -----------------------------------------------------------

    /// Bun's `Cmd` write_failing_error: a failed write rejects the run.
    pub(crate) async fn cmd_write_failing_error(
        &self,
        io: &ShellIO,
        shell: &ShellExecEnv,
        msg: &str,
    ) -> ExecResult {
        match write_err(io, shell, msg).await {
            Ok(()) => Ok(1),
            Err(e) => Err(ShellError::system(e.message())),
        }
    }

    /// The CondExpr/Assigns flavour: a failed write becomes the exit code.
    async fn write_failing_error_no_throw(
        &self,
        io: &ShellIO,
        shell: &ShellExecEnv,
        msg: &str,
    ) -> i32 {
        match write_err(io, shell, msg).await {
            Ok(()) => 1,
            Err(e) => e.errno,
        }
    }
}

async fn write_err(io: &ShellIO, shell: &ShellExecEnv, msg: &str) -> Result<(), ShellSysError> {
    io.stderr
        .write(msg.as_bytes(), &shell.buffered_stderr)
        .await
}

#[cfg(unix)]
fn is_char_device(st: &std::fs::Metadata) -> bool {
    use std::os::unix::fs::FileTypeExt;
    st.file_type().is_char_device()
}

#[cfg(not(unix))]
fn is_char_device(_st: &std::fs::Metadata) -> bool {
    false
}

#[cfg(test)]
mod tests;