diff --git a/src/runtime/shell/states/Pipeline.rs b/src/runtime/shell/states/Pipeline.rs index 3e835b0c39f7..8e15a78e31dd 100644 --- a/src/runtime/shell/states/Pipeline.rs +++ b/src/runtime/shell/states/Pipeline.rs @@ -18,8 +18,8 @@ pub struct Pipeline { pub node: bun_ptr::BackRef, pub(crate) io: IO, pub(crate) exited_count: u32, + /// `None` until `setup_commands` has inited every child. pub(crate) cmds: Option>, - pub(crate) pipes: Option>, pub(crate) state: PipelineState, } @@ -29,10 +29,15 @@ pub enum CmdOrResult { } pub enum PipelineState { - StartingCmds { idx: u32 }, + /// `idx` is the next `cmds[]` slot to start. + StartingCmds { + idx: u32, + }, Pending, WaitingWriteErr, - Done { exit_code: ExitCode }, + Done { + exit_code: ExitCode, + }, } impl Default for PipelineState { @@ -55,7 +60,6 @@ impl Pipeline { io, exited_count: 0, cmds: None, - pipes: None, state: PipelineState::default(), })) } @@ -89,24 +93,44 @@ impl Pipeline { } } - /// Set up N-1 pipes, dupe the shell env per child, spawn each - /// Cmd/Assigns/Subshell/If/CondExpr with stdin/stdout wired to the right - /// pipe ends. - /// - /// Spawns exactly ONE child - /// per call and returns that child's `start()` Yield. The trampoline's - /// `drain_pipelines` (Yield.rs) re-enters `Pipeline::next` to spawn the - /// next child once the current one suspends — so every child's start-yield - /// is driven, never dropped. + /// Starts ONE child per call; `drain_pipelines` (Yield.rs) re-enters + /// `Pipeline::next` for the next one once the current child suspends. fn next_starting(interp: &Interpreter, this: NodeId, idx: u32) -> Yield { + if interp.as_pipeline(this).cmds.is_none() { + debug_assert_eq!(idx, 0); + if let Some(y) = Self::setup_commands(interp, this) { + return y; + } + } + + let next = { + let me = interp.as_pipeline(this); + let cmds = me.cmds.as_deref().expect("set by setup_commands"); + cmds.get(idx as usize).map(|slot| match slot { + CmdOrResult::Cmd(id) => *id, + CmdOrResult::Result(_) => { + unreachable!("pipeline child {} finished before it was started", idx) + } + }) + }; + let Some(child) = next else { + // All children started; wait for their `child_done` callbacks. + interp.as_pipeline_mut(this).state = PipelineState::Pending; + return Yield::suspended(); + }; + interp.as_pipeline_mut(this).state = PipelineState::StartingCmds { idx: idx + 1 }; + interp.start_node(child) + } + + /// Creates the pipes and inits every child without starting any, so a + /// failed pipe or dup finishes the pipeline (`Some(yield)`) while no child + /// runs a subtree that `deinit` cannot reach. + fn setup_commands(interp: &Interpreter, this: NodeId) -> Option { let (node, parent_shell, evtloop) = { let me = interp.as_pipeline(this); (me.node, me.base.shell, interp.event_loop) }; let items: &[ast::PipelineItem] = node.items; - // Assigns inside a pipeline are - // no-ops — they're not counted, not duped, not started. `cmd_count` - // here is the number of *runnable* children. let cmd_count = items .iter() .filter(|it| !matches!(it, ast::PipelineItem::Assigns(_))) @@ -114,151 +138,121 @@ impl Pipeline { if cmd_count == 0 { // An empty pipeline finishes with 0. - return Self::finish(interp, this, 0); + return Some(Self::finish(interp, this, 0)); } - // First entry: allocate pipes + cmd slots. - if idx == 0 && interp.as_pipeline(this).cmds.is_none() { - let mut pipes: Vec = Vec::with_capacity(cmd_count.saturating_sub(1)); - for _ in 0..cmd_count.saturating_sub(1) { - // On POSIX use a - // UNIX stream socketpair via `socketpairForShell` — on macOS - // that variant intentionally skips SO_NOSIGPIPE so the - // subprocess writing to a closed read end is killed by SIGPIPE - // (like a real shell) instead of seeing EPIPE and printing - // "Broken pipe" to stderr; on Windows use pipe(). - #[cfg(windows)] - let r = bun_sys::pipe(); - #[cfg(unix)] - let r = bun_sys::socketpair_for_shell(libc::AF_UNIX, libc::SOCK_STREAM, 0, false); - match r { - Ok(p) => pipes.push(p), - Err(e) => { - for p in &pipes { - closefd(p[0]); - closefd(p[1]); - } - let sys_err = e.to_shell_system_error(); - return Self::write_failing_error( - interp, - this, - format_args!("bun: {}\n", sys_err.message), - ); + let mut pipes: Vec = Vec::with_capacity(cmd_count - 1); + for _ in 0..cmd_count - 1 { + // On POSIX use a + // UNIX stream socketpair via `socketpairForShell` — on macOS + // that variant intentionally skips SO_NOSIGPIPE so the + // subprocess writing to a closed read end is killed by SIGPIPE + // (like a real shell) instead of seeing EPIPE and printing + // "Broken pipe" to stderr; on Windows use pipe(). + #[cfg(windows)] + let r = bun_sys::pipe(); + #[cfg(unix)] + let r = bun_sys::socketpair_for_shell(libc::AF_UNIX, libc::SOCK_STREAM, 0, false); + match r { + Ok(p) => pipes.push(p), + Err(e) => { + for p in &pipes { + closefd(p[0]); + closefd(p[1]); } + let sys_err = e.to_shell_system_error(); + return Some(Self::write_failing_error( + interp, + this, + format_args!("bun: {}\n", sys_err.message), + )); } } - let cmds: Vec = (0..cmd_count).map(|_| CmdOrResult::Result(0)).collect(); - let me = interp.as_pipeline_mut(this); - me.pipes = Some(pipes.into_boxed_slice()); - me.cmds = Some(cmds.into_boxed_slice()); } - // `idx` walks `items[]`; skip over Assigns to find the next runnable. - let mut item_idx = idx as usize; - while item_idx < items.len() && matches!(items[item_idx], ast::PipelineItem::Assigns(_)) { - item_idx += 1; - } - if item_idx >= items.len() { - // All children spawned; wait for their `child_done` callbacks. - interp.as_pipeline_mut(this).state = PipelineState::Pending; - return Yield::suspended(); - } - // `cmd_idx` is the position among runnable children (indexes - // `pipes[]`/`cmds[]`). - let cmd_idx = items[..item_idx] - .iter() - .filter(|it| !matches!(it, ast::PipelineItem::Assigns(_))) - .count(); - - // Build per-child IO: stdin from prev pipe read end (or parent - // stdin for first), stdout to this pipe write end (or parent stdout - // for last), stderr inherited. let interp_ptr: *mut Interpreter = interp.as_ctx_ptr(); - let child_io = { - let me = interp.as_pipeline(this); - let pipes = me.pipes.as_ref().expect("pipes set above"); - let stdin = if cmd_count == 1 || cmd_idx == 0 { - me.io.stdin.clone() - } else { - let r = IOReader::init(pipes[cmd_idx - 1][0], evtloop); - r.set_interp(interp_ptr); - InKind::Fd(r) - }; - let stdout = if cmd_count == 1 || cmd_idx == cmd_count - 1 { - me.io.stdout.clone() - } else { - // `is_socket` is set on POSIX — the POSIX - // pipe is actually a socketpair end (see above). - let w = IOWriter::init( - pipes[cmd_idx][1], - io_writer::Flags { - pollable: true, - is_socket: cfg!(unix), - ..Default::default() - }, - evtloop, - ); - w.set_interp(interp_ptr); - OutKind::Fd(crate::shell::io::OutFd { - writer: w, - captured: None, - }) - }; - IO { - stdin, - stdout, - stderr: me.io.stderr.clone(), + let mut cmds: Vec = Vec::with_capacity(cmd_count); + for item in items { + if matches!(item, ast::PipelineItem::Assigns(_)) { + continue; } - }; + // Position among runnable children (indexes `pipes[]`/`cmds[]`). + let cmd_idx = cmds.len(); - // Each pipeline child gets its own duped env (var assignments - // inside a pipeline must not leak to siblings or the parent). - // SAFETY: `parent_shell` is a live env owned by this pipeline's - // parent state. - let duped = match unsafe { - (*parent_shell).dupe_for_subshell(&child_io, ShellExecEnvKind::Pipeline) - } { - Ok(d) => d, - Err(e) => { - // Drop `child_io` (its IOReader/IOWriter own two of the pipe - // ends) and close the pipe ends no child has claimed yet. - drop(child_io); - { - let me = interp.as_pipeline_mut(this); - if let Some(pipes) = me.pipes.as_ref() { - let len = pipes.len(); - for p in &pipes[cmd_idx..] { - closefd(p[0]); - } - for p in &pipes[core::cmp::min(cmd_idx + 1, len)..] { - closefd(p[1]); - } - } + let child_io = { + let me = interp.as_pipeline(this); + let stdin = if cmd_idx == 0 { + me.io.stdin.clone() + } else { + let r = IOReader::init(pipes[cmd_idx - 1][0], evtloop); + r.set_interp(interp_ptr); + InKind::Fd(r) + }; + let stdout = if cmd_idx == cmd_count - 1 { + me.io.stdout.clone() + } else { + // `is_socket` is set on POSIX — the POSIX + // pipe is actually a socketpair end (see above). + let w = IOWriter::init( + pipes[cmd_idx][1], + io_writer::Flags { + pollable: true, + is_socket: cfg!(unix), + ..Default::default() + }, + evtloop, + ); + w.set_interp(interp_ptr); + OutKind::Fd(crate::shell::io::OutFd { + writer: w, + captured: None, + }) + }; + IO { + stdin, + stdout, + stderr: me.io.stderr.clone(), } - let sys_err = e.to_shell_system_error(); - return Self::write_failing_error( - interp, - this, - format_args!("bun: {}\n", sys_err.message), - ); - } - }; + }; - let child = match items[item_idx] { - ast::PipelineItem::Cmd(c) => Cmd::init(interp, duped, c, this, child_io), - ast::PipelineItem::Subshell(s) => Subshell::init(interp, duped, s, this, child_io), - ast::PipelineItem::If(f) => If::init(interp, duped, f, this, child_io), - ast::PipelineItem::CondExpr(c) => CondExpr::init(interp, duped, c, this, child_io), - ast::PipelineItem::Assigns(_) => unreachable!("skipped above"), - }; - interp.as_pipeline_mut(this).cmds.as_mut().unwrap()[cmd_idx] = CmdOrResult::Cmd(child); - interp.as_pipeline_mut(this).state = PipelineState::StartingCmds { - idx: (item_idx + 1) as u32, - }; + // Each pipeline child gets its own duped env (var assignments + // inside a pipeline must not leak to siblings or the parent). + // SAFETY: `parent_shell` is a live env owned by this pipeline's + // parent state. + let duped = match unsafe { + (*parent_shell).dupe_for_subshell(&child_io, ShellExecEnvKind::Pipeline) + } { + Ok(d) => d, + Err(e) => { + // Close the pipe ends that neither `child_io` nor an inited child owns. + drop(child_io); + for p in &pipes[cmd_idx..] { + closefd(p[0]); + } + for p in &pipes[core::cmp::min(cmd_idx + 1, pipes.len())..] { + closefd(p[1]); + } + interp.as_pipeline_mut(this).cmds = Some(cmds.into_boxed_slice()); + let sys_err = e.to_shell_system_error(); + return Some(Self::write_failing_error( + interp, + this, + format_args!("bun: {}\n", sys_err.message), + )); + } + }; - // Spawn exactly this one child. The trampoline will re-enter us via - // `drain_pipelines` to spawn the next after this one yields. - interp.start_node(child) + let child = match *item { + ast::PipelineItem::Cmd(c) => Cmd::init(interp, duped, c, this, child_io), + ast::PipelineItem::Subshell(s) => Subshell::init(interp, duped, s, this, child_io), + ast::PipelineItem::If(f) => If::init(interp, duped, f, this, child_io), + ast::PipelineItem::CondExpr(c) => CondExpr::init(interp, duped, c, this, child_io), + ast::PipelineItem::Assigns(_) => unreachable!("skipped above"), + }; + cmds.push(CmdOrResult::Cmd(child)); + } + interp.as_pipeline_mut(this).cmds = Some(cmds.into_boxed_slice()); + None } /// Mark the pipeline done with `exit_code`. Returns `Next(this)` so the @@ -282,15 +276,15 @@ impl Pipeline { use std::io::Write as _; let mut buf = Vec::new(); let _ = buf.write_fmt(args); - if interp.as_pipeline(this).io.stderr.needs_io().is_some() { + let stderr_fd = match &interp.as_pipeline(this).io.stderr { + OutKind::Fd(fd) => Some((std::sync::Arc::clone(&fd.writer), fd.captured)), + OutKind::Pipe | OutKind::Ignore => None, + }; + if let Some((writer, captured)) = stderr_fd { // Only the fd arm transitions state. interp.as_pipeline_mut(this).state = PipelineState::WaitingWriteErr; let child = io_writer::ChildPtr::new(this, io_writer::WriterTag::Pipeline); - // `OutKind::Fd` guaranteed by `needs_io()`. - if let OutKind::Fd(fd) = &interp.as_pipeline(this).io.stderr { - return fd.writer.enqueue(child, fd.captured, &buf); - } - unreachable!() + return writer.enqueue(child, captured, &buf); } if let OutKind::Pipe = &interp.as_pipeline(this).io.stderr { // SAFETY: single trampoline frame; no other borrow of the env's @@ -398,7 +392,7 @@ impl Pipeline { pub(crate) fn deinit(interp: &Interpreter, this: NodeId) { log!("Pipeline {} deinit", this); - // Deinit any still-live children (and their duped envs). + // Only children that never started (setup failed) are still here. let cmds = interp.as_pipeline_mut(this).cmds.take(); if let Some(cmds) = cmds { for c in cmds.into_vec() { @@ -408,10 +402,5 @@ impl Pipeline { } } } - let me = interp.as_pipeline_mut(this); - // The pipe fds are owned by the IOReader/IOWriter Arcs handed to each - // child; when those drop they close. Any unclaimed ones (error path) - // were closed inline above. - me.pipes = None; } } diff --git a/test/js/bun/shell/bunshell.test.ts b/test/js/bun/shell/bunshell.test.ts index 2f75fcad3f5d..d012d53c6e40 100644 --- a/test/js/bun/shell/bunshell.test.ts +++ b/test/js/bun/shell/bunshell.test.ts @@ -647,6 +647,38 @@ describe("bunshell", () => { expect(results).toEqual([...Array(firstFailure).fill(ok), ...Array(results.length - firstFailure).fill(emfile)]); expect(exitCode).toBe(0); }); + + // Same sweep with a subshell at the head. A subshell runs a Script of its + // own, so if the pipeline started it before a later command's `dup` + // failed, tearing the pipeline down would leave that Script behind with + // the subshell's echo still writing into the pipe. The echo then reports + // to a freed node once the pipe is closed under it. No child may start + // until every child is set up. + test("reports EMFILE from the per-command env dup when the head is a subshell", async () => { + const script = /* ts */ ` + import { $ } from "bun"; + const results = []; + for (let n = 2; n <= 16; n++) { + const pipeline = ["(echo hi)", ...Array(n - 1).fill("cat")].join(" | "); + const r = await $\`\${{ raw: pipeline }}\`.nothrow().quiet(); + results.push({ stdout: r.stdout.toString(), stderr: r.stderr.toString(), exitCode: r.exitCode }); + // One event loop turn, so that a write a torn-down child left in + // flight completes (and reports) before the next pipeline runs. + await Bun.sleep(0); + } + console.log(JSON.stringify(results)); + `; + await using proc = runWithFdLimit(32, script, { ...bunEnv, BUN_ENABLE_EXPERIMENTAL_SHELL_BUILTINS: "1" }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + const results = JSON.parse(stdout); + const ok = { stdout: "hi\n", stderr: "", exitCode: 0 }; + const emfile = { stdout: "", stderr: "bun: Too many open files\n", exitCode: 1 }; + const firstFailure = results.findIndex(r => r.exitCode !== 0); + expect(firstFailure).toBeGreaterThan(0); + expect(results).toEqual([...Array(firstFailure).fill(ok), ...Array(results.length - firstFailure).fill(emfile)]); + expect(exitCode).toBe(0); + }); }); describe("operators no spaces", async () => {