diff --git a/src/runtime/shell/IOReader.rs b/src/runtime/shell/IOReader.rs index 3c02f7ce4759..c39ec12171c8 100644 --- a/src/runtime/shell/IOReader.rs +++ b/src/runtime/shell/IOReader.rs @@ -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). @@ -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, @@ -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(); + } + self.state().is_reading = true; if let Err(e) = self.reader().start_with_current_pipe() { self.on_reader_error(&e); return Yield::failed(); @@ -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 = 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) { @@ -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 = 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) { diff --git a/test/js/bun/shell/bunshell.test.ts b/test/js/bun/shell/bunshell.test.ts index d1e6553875e7..8828251848ce 100644 --- a/test/js/bun/shell/bunshell.test.ts +++ b/test/js/bun/shell/bunshell.test.ts @@ -1502,6 +1502,16 @@ describe("deno_task", () => { .stdout("0\n") .runAsTest("long pipeline"); + // Every `cat` here shares the subshell's stdin Arc. 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")