Skip to content
Open
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
50 changes: 29 additions & 21 deletions src/runtime/shell/IOReader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,8 @@ struct State {
evtloop: EventLoopHandle,
#[cfg(windows)]
is_reading: bool,
/// Set while `drain_readers` runs; a nested call returns and the outer loop handles new entries.
draining: bool,
/// Weak self-ref so `keepalive()` can bump the strong count from `&self`
/// without unsafe Arc-pointer reconstruction. Set via `Arc::new_cyclic` in
/// `init()` (the sole constructor).
Expand Down Expand Up @@ -141,6 +143,7 @@ impl IOReader {
evtloop,
#[cfg(windows)]
is_reading: false,
draining: false,
self_weak: std::sync::Weak::clone(w),
read_guards: Vec::new(),
interp: None,
Expand Down Expand Up @@ -224,11 +227,15 @@ impl IOReader {
}
#[cfg(windows)]
{
let s = self.state();
if s.is_reading {
if self.state().is_reading {
return Yield::suspended();
}
s.is_reading = true;
// Already at EOF (the source is gone): just notify whoever registered late.
if self.reader().is_done() {
self.drain_readers();
return Yield::suspended();
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
self.state().is_reading = true;
if let Err(e) = self.reader().start_with_current_pipe() {
self.on_reader_error(&e);
return Yield::failed();
Expand Down Expand Up @@ -306,17 +313,8 @@ impl IOReader {
// alive across the loop.
let _keepalive = self.keepalive();
self.set_reading(false);
let s = self.state();
s.raw_err = Some(err.clone());
// NOTE: reshaped for borrowck — copy out before dispatching.
let readers: Vec<ChildPtr> = s.readers.clone();
let interp = s.interp;
for r in readers {
// Re-derive a fresh SystemError per callee (see
// IOWriter.on_error note).
let ee = err.to_shell_system_error();
self.run_yield(dispatch_reader_done(r, Some(ee), interp));
}
self.state().raw_err = Some(err.clone());
self.drain_readers();
}

fn on_reader_done_cb(&self) {
Expand All @@ -326,17 +324,27 @@ impl IOReader {
// Hold a strong ref across the body.
let _keepalive = self.keepalive();
self.set_reading(false);
self.drain_readers();
}

/// Pops before dispatching: a callback may `add_reader` (must be notified too) or free this entry's node.
fn drain_readers(&self) {
let s = self.state();
let readers: Vec<ChildPtr> = s.readers.clone();
if s.draining {
return;
}
s.draining = true;
let interp = s.interp;
// `SystemError` isn't `Clone` yet, so we keep the source `sys::Error`
// (which IS `Clone`) and re-derive a fresh `SystemError` per callee —
// same approach as `on_reader_error`.
let raw_err = s.raw_err.clone();
for r in readers {
let ee = raw_err.as_ref().map(|e| e.to_shell_system_error());
while !self.state().readers.is_empty() {
let r = self.state().readers.swap_remove(0);
let ee = self
.state()
.raw_err
.as_ref()
.map(|e| e.to_shell_system_error());
self.run_yield(dispatch_reader_done(r, ee, interp));
}
self.state().draining = false;
}

fn run_yield(&self, y: Yield) {
Expand Down
10 changes: 10 additions & 0 deletions test/js/bun/shell/bunshell.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1502,6 +1502,16 @@ describe("deno_task", () => {
.stdout("0\n")
.runAsTest("long pipeline");

// Every `cat` here shares the subshell's stdin Arc<IOReader>. When EOF fires, each
// reader's done-handler starts the next `cat` via the Yield trampoline, which calls
// add_reader()+start() on the same already-done IOReader. drain_readers() must pop
// each entry before dispatching (so the readers Vec can mutate safely and add_reader's
// dedup never matches a freed-then-reused NodeId) and start() must drain
// late-registered readers rather than restarting the finished pipe.
TestBuilder.command`echo hi | (cat && cat && cat && cat && cat && cat && cat && cat && cat && cat && cat && cat)`
.stdout("hi\n")
.runAsTest("many readers on shared stdin IOReader");

// Test pipeline stack consistency with complex nesting
TestBuilder.command`echo outer | (echo inner1 | echo inner2 | (echo deep1 | echo deep2) | echo inner3) | echo final | BUN_TEST_VAR=1 ${BUN} -e 'process.stdin.pipe(process.stdout)'`
.stdout("final\n")
Expand Down
Loading