Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion src/runtime/shell/IOWriter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ pub enum WriterTag {
Builtin,
Cmd,
CondExpr,
Pipeline,
/// `subproc::PipeReader::CapturedWriter` — heap-allocated, addressed via
/// `ChildPtr::raw` rather than `node`.
Subproc,
Expand Down Expand Up @@ -1225,13 +1226,16 @@ pub(crate) fn on_io_writer_chunk(
err: Option<sys::SystemError>,
) -> Yield {
use crate::shell::builtin::Builtin;
use crate::shell::states::{cmd, cond_expr};
use crate::shell::states::{cmd, cond_expr, pipeline};
match child.tag {
WriterTag::Builtin => Builtin::on_io_writer_chunk(interp, child.node, written, err),
WriterTag::Cmd => cmd::Cmd::on_io_writer_chunk(interp, child.node, written, err),
WriterTag::CondExpr => {
cond_expr::CondExpr::on_io_writer_chunk(interp, child.node, written, err)
}
WriterTag::Pipeline => {
pipeline::Pipeline::on_io_writer_chunk(interp, child.node, written, err)
}
// The target is the subprocess PipeReader's `CapturedWriter`; it
// lives outside the NodeId arena (heap-allocated PipeReader), so it
// is carried in `child.raw` instead of `child.node`.
Expand Down
102 changes: 79 additions & 23 deletions src/runtime/shell/states/Pipeline.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
use crate::shell::ExitCode;
use crate::shell::ast;
use crate::shell::interpreter::{
Interpreter, Node, NodeId, Pipe, ShellExecEnv, ShellExecEnvKind, StateKind, closefd, log,
Expand All @@ -11,7 +12,6 @@ use crate::shell::states::cond_expr::CondExpr;
use crate::shell::states::r#if::If;
use crate::shell::states::subshell::Subshell;
use crate::shell::yield_::Yield;
use crate::shell::{ExitCode, ShellErr};

pub struct Pipeline {
pub(crate) base: Base,
Expand Down Expand Up @@ -114,12 +114,7 @@ impl Pipeline {

if cmd_count == 0 {
// An empty pipeline finishes with 0.
// Return `Next(this)` so the trampoline sees `is_done`, removes us
// from the pipeline stack, and `next()` bubbles to the parent.
// Calling `child_done(parent, ..)` directly here would free this
// node while it's still on `pipeline_stack`.
interp.as_pipeline_mut(this).state = PipelineState::Done { exit_code: 0 };
return Yield::Next(this);
return Self::finish(interp, this, 0);
}

// First entry: allocate pipes + cmd slots.
Expand All @@ -143,13 +138,12 @@ impl Pipeline {
closefd(p[0]);
closefd(p[1]);
}
// Leave `StartingCmds` so `drain_pipelines` doesn't
// re-enter `next_starting` and retry the failing
// syscall in a loop; stay suspended until the error
// write completes.
interp.as_pipeline_mut(this).state = PipelineState::WaitingWriteErr;
interp.throw(ShellErr::new_sys(&e));
return Yield::failed();
let sys_err = e.to_shell_system_error();
return Self::write_failing_error(
interp,
this,
format_args!("bun: {}\n", sys_err.message),
);
}
}
}
Expand Down Expand Up @@ -226,12 +220,8 @@ impl Pipeline {
} {
Ok(d) => d,
Err(e) => {
// On dupe failure,
// close the pipe ends not yet wrapped in an IOReader/IOWriter,
// deref `cmd_io`, transition to `.waiting_write_err`, and
// suspend. Without the state transition `drain_pipelines`
// would re-enter at the same `idx`, re-wrapping the same fds
// in fresh IOReader/IOWriter each iteration.
// 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);
Expand All @@ -244,10 +234,13 @@ impl Pipeline {
closefd(p[1]);
}
}
me.state = PipelineState::WaitingWriteErr;
}
interp.throw(ShellErr::new_sys(&e));
return Yield::failed();
let sys_err = e.to_shell_system_error();
return Self::write_failing_error(
interp,
this,
format_args!("bun: {}\n", sys_err.message),
);
}
};

Expand All @@ -268,6 +261,69 @@ impl Pipeline {
interp.start_node(child)
}

/// Mark the pipeline done with `exit_code`. Returns `Next(this)` so the
/// trampoline sees `is_done`, removes us from `pipeline_stack`, and
/// `next()` reports to the parent. Calling `child_done(parent, ..)`
/// directly would free this node while it is still on `pipeline_stack`.
fn finish(interp: &Interpreter, this: NodeId, exit_code: ExitCode) -> Yield {
interp.as_pipeline_mut(this).state = PipelineState::Done { exit_code };
Yield::Next(this)
}

/// Same shape as `Builtin::cmd_write_failing_error`: `.fd` stderr
/// enqueues an async write and parks in `WaitingWriteErr` (resumed by
/// `on_io_writer_chunk`); otherwise append to the captured stderr buffer
/// and finish with exit 1.
fn write_failing_error(
interp: &Interpreter,
this: NodeId,
args: core::fmt::Arguments<'_>,
) -> Yield {
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() {
// 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!()
}
if let OutKind::Pipe = &interp.as_pipeline(this).io.stderr {
// SAFETY: single trampoline frame; no other borrow of the env's
// (or its parent's) stderr buffer is live.
let stderr = unsafe {
interp
.as_pipeline_mut(this)
.base
.shell_mut()
.buffered_stderr_mut()
};
stderr.extend_from_slice(&buf);
}
Self::finish(interp, this, 1)
}

/// IOWriter completion callback for the error message written in
/// `WaitingWriteErr`. The pipeline finishes with exit code 1 whether or
/// not the write succeeded: the parent always needs a completion, and a
/// failed stderr write has nowhere else to be reported.
pub(crate) fn on_io_writer_chunk(
interp: &Interpreter,
this: NodeId,
_written: usize,
_err: Option<bun_sys::SystemError>,
) -> Yield {
debug_assert!(matches!(
interp.as_pipeline(this).state,
PipelineState::WaitingWriteErr
));
Self::finish(interp, this, 1)
}

pub(crate) fn child_done(
interp: &Interpreter,
this: NodeId,
Expand Down
75 changes: 75 additions & 0 deletions test/js/bun/shell/bunshell.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -574,6 +574,81 @@ describe("bunshell", () => {
expect(exitCode).toBe(0);
});

// `ulimit -n` caps the fd table of the child so the pipeline cannot create
// its pipes (EMFILE). The pipeline must print the error on its stderr and
// finish with exit code 1 so the `$` promise settles and the script goes on,
// instead of throwing from inside the interpreter and never completing.
describe.skipIf(isWindows)("pipeline that fails to create its pipes", () => {
const runWithFdLimit = (limit: number, script: string, env = bunEnv) =>
Bun.spawn({
cmd: ["/bin/sh", "-c", `ulimit -n ${limit} && exec "$1" -e "$2"`, "sh", bunExe(), script],
env,
stdout: "pipe",
stderr: "pipe",
});

test("reports EMFILE on stderr and finishes with exit code 1", async () => {
// 24 pipes need 48 fds; the whole process gets 32.
const pipeline = ["echo hi", ...Array(24).fill("cat")].join(" | ");
const script = /* ts */ `
import { $ } from "bun";
const results = {};
const record = (name, r) => {
results[name] = { stdout: r.stdout.toString(), stderr: r.stderr.toString(), exitCode: r.exitCode };
};
record("stderr captured", await $\`${pipeline}\`.nothrow().quiet());
record("stderr on a fd", await $\`${pipeline}\`.nothrow());
record("script continues", await $\`${pipeline}; echo after\`.nothrow().quiet());
try {
await $\`${pipeline}\`.quiet();
results["throws"] = "did not throw";
} catch (e) {
record("throws", e);
}
console.log(JSON.stringify(results));
`;
await using proc = runWithFdLimit(32, script);
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(stderr).toBe("bun: Too many open files\n");
expect(JSON.parse(stdout)).toEqual({
"stderr captured": { stdout: "", stderr: "bun: Too many open files\n", exitCode: 1 },
"stderr on a fd": { stdout: "", stderr: "bun: Too many open files\n", exitCode: 1 },
"script continues": { stdout: "after\n", stderr: "bun: Too many open files\n", exitCode: 0 },
"throws": { stdout: "", stderr: "bun: Too many open files\n", exitCode: 1 },
});
expect(exitCode).toBe(0);
});

// Pipelines of growing length against a fixed fd budget. The short ones
// fit. Past some length the pipes still fit but the per-command `dup` of
// the cwd fd fails. Past that the pipes themselves fail. Every pipeline
// must finish, either with its output or with the EMFILE message.
test("reports EMFILE from the per-command env dup and finishes", 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 });
}
console.log(JSON.stringify(results));
`;
// The builtin cat keeps this to pipes and dups: no subprocess, so no
// spawn-time fds change the accounting.
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 () => {
TestBuilder.command`echo LMAO|cat`.stdout("LMAO\n").runAsTest("pipeline");
TestBuilder.command`echo foo&&echo hi`.stdout("foo\nhi\n").runAsTest("&&");
Expand Down
Loading