From 95dfb3981c4c17bb82a965788f1ad2ec83241d5a Mon Sep 17 00:00:00 2001 From: naokihaba <59875779+naokihaba@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:57:23 +0900 Subject: [PATCH 1/5] feat(retrying_writer): implement a writer that retries on WouldBlock errors --- .../src/session/reporter/retrying_writer.rs | 266 ++++++++++++++++++ 1 file changed, 266 insertions(+) create mode 100644 crates/vt/src/session/reporter/retrying_writer.rs diff --git a/crates/vt/src/session/reporter/retrying_writer.rs b/crates/vt/src/session/reporter/retrying_writer.rs new file mode 100644 index 000000000..0fd79448f --- /dev/null +++ b/crates/vt/src/session/reporter/retrying_writer.rs @@ -0,0 +1,266 @@ +//! A [`Write`] wrapper that waits and retries while its sink is temporarily full. + +use std::io::{self, Write}; + +/// A sink that can wait until it accepts more bytes. +pub trait WaitWritable { + /// Returns once a write to the sink is worth retrying. + fn wait_writable(&self) -> io::Result<()>; +} + +/// Writer that retries a write or flush that fails with [`io::ErrorKind::WouldBlock`]. +/// +/// `vp run` shares its stdout and stderr with the tasks that inherit them. +/// Such a task can switch the shared open file description to non-blocking +/// mode (Node.js does this when the stream is a pipe), and a write to a full +/// pipe then fails with `EAGAIN` instead of waiting. +/// +/// The retry has to sit directly above the stream. A writer further up, such +/// as `LabeledWriter`, or a caller of `write_all` can't tell how many bytes a +/// failed call wrote, so it can't resume without repeating or dropping some. +pub struct RetryingWriter { + inner: W, +} + +impl RetryingWriter { + pub const fn new(inner: W) -> Self { + Self { inner } + } +} + +impl Write for RetryingWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + loop { + match self.inner.write(buf) { + Err(err) if err.kind() == io::ErrorKind::WouldBlock => { + self.inner.wait_writable()?; + } + result => return result, + } + } + } + + fn flush(&mut self) -> io::Result<()> { + loop { + match self.inner.flush() { + Err(err) if err.kind() == io::ErrorKind::WouldBlock => { + self.inner.wait_writable()?; + } + result => return result, + } + } + } +} + +/// Sleeps until `fd` accepts a write, or until it reports an error or hangup, +/// which the retried write then returns. +#[cfg(unix)] +fn wait_fd_writable(fd: std::os::fd::BorrowedFd<'_>) -> io::Result<()> { + use nix::{ + errno::Errno, + poll::{PollFd, PollFlags, PollTimeout, poll}, + }; + + let mut fds = [PollFd::new(fd, PollFlags::POLLOUT)]; + match poll(&mut fds, PollTimeout::NONE) { + Ok(_) | Err(Errno::EINTR) => Ok(()), + Err(errno) => Err(errno.into()), + } +} + +macro_rules! impl_wait_writable_for_std_stream { + ($stream:ty) => { + impl WaitWritable for $stream { + #[cfg(unix)] + fn wait_writable(&self) -> io::Result<()> { + use std::os::fd::AsFd as _; + + wait_fd_writable(self.as_fd()) + } + + // Windows has no shared non-blocking mode for a task to switch on, + // so a `WouldBlock` isn't expected here. Yield and retry if one + // appears anyway. + #[cfg(windows)] + fn wait_writable(&self) -> io::Result<()> { + std::thread::yield_now(); + Ok(()) + } + } + }; +} + +impl_wait_writable_for_std_stream!(io::Stdout); +impl_wait_writable_for_std_stream!(io::Stderr); + +#[cfg(test)] +mod tests { + use std::cell::Cell; + + use super::*; + + /// A sink like a non-blocking pipe: it takes `capacity` bytes, then + /// returns `WouldBlock` until `wait_writable` empties it. + struct FullSink { + capacity: usize, + free: Cell, + waits: Cell, + flush_blocks: usize, + written: Vec, + } + + impl FullSink { + fn new(capacity: usize) -> Self { + Self { + capacity, + free: Cell::new(capacity), + waits: Cell::new(0), + flush_blocks: 0, + written: Vec::new(), + } + } + } + + impl Write for FullSink { + fn write(&mut self, buf: &[u8]) -> io::Result { + let n = buf.len().min(self.free.get()); + if n == 0 { + return Err(io::ErrorKind::WouldBlock.into()); + } + self.free.set(self.free.get() - n); + self.written.extend_from_slice(&buf[..n]); + Ok(n) + } + + fn flush(&mut self) -> io::Result<()> { + if self.flush_blocks == 0 { + return Ok(()); + } + self.flush_blocks -= 1; + Err(io::ErrorKind::WouldBlock.into()) + } + } + + impl WaitWritable for FullSink { + fn wait_writable(&self) -> io::Result<()> { + self.waits.set(self.waits.get() + 1); + self.free.set(self.capacity); + Ok(()) + } + } + + #[test] + fn write_all_waits_while_the_sink_is_full() { + let mut writer = RetryingWriter::new(FullSink::new(3)); + writer.write_all(b"0123456789").unwrap(); + + assert_eq!(writer.inner.written, b"0123456789"); + assert_eq!(writer.inner.waits.get(), 3); + } + + #[test] + fn flush_waits_while_the_sink_is_full() { + let mut sink = FullSink::new(3); + sink.flush_blocks = 2; + let mut writer = RetryingWriter::new(sink); + writer.flush().unwrap(); + + assert_eq!(writer.inner.waits.get(), 2); + } + + #[test] + fn other_errors_are_returned() { + struct Closed; + + impl Write for Closed { + fn write(&mut self, _buf: &[u8]) -> io::Result { + Err(io::ErrorKind::BrokenPipe.into()) + } + + fn flush(&mut self) -> io::Result<()> { + Err(io::ErrorKind::BrokenPipe.into()) + } + } + + impl WaitWritable for Closed { + fn wait_writable(&self) -> io::Result<()> { + panic!("a broken pipe must not be retried"); + } + } + + let mut writer = RetryingWriter::new(Closed); + assert_eq!(writer.write(b"x").unwrap_err().kind(), io::ErrorKind::BrokenPipe); + assert_eq!(writer.flush().unwrap_err().kind(), io::ErrorKind::BrokenPipe); + } + + /// The real case: a pipe in non-blocking mode that is full when the write + /// starts. The reader starts only after the writer has begun to wait, so + /// the write can't succeed without the retry. + #[cfg(unix)] + #[test] + fn write_all_waits_for_a_full_non_blocking_pipe() { + use std::{ + io::Read as _, + os::fd::AsFd as _, + sync::mpsc::{Sender, channel}, + }; + + use nix::fcntl::{FcntlArg, OFlag, fcntl}; + + struct PipeSink { + pipe: io::PipeWriter, + waiting: Sender<()>, + } + + impl Write for PipeSink { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.pipe.write(buf) + } + + fn flush(&mut self) -> io::Result<()> { + self.pipe.flush() + } + } + + impl WaitWritable for PipeSink { + fn wait_writable(&self) -> io::Result<()> { + // The reader may already be gone; it only needs the first signal. + let _ = self.waiting.send(()); + wait_fd_writable(self.pipe.as_fd()) + } + } + + let (mut reader, mut pipe) = io::pipe().unwrap(); + let flags = OFlag::from_bits_retain(fcntl(pipe.as_fd(), FcntlArg::F_GETFL).unwrap()); + fcntl(pipe.as_fd(), FcntlArg::F_SETFL(flags | OFlag::O_NONBLOCK)).unwrap(); + + let mut filled = 0; + loop { + match pipe.write(&[b'.'; 4096]) { + Ok(n) => filled += n, + Err(err) => { + assert_eq!(err.kind(), io::ErrorKind::WouldBlock); + break; + } + } + } + + let (waiting, started_waiting) = channel(); + let reader_thread = std::thread::spawn(move || { + started_waiting.recv().unwrap(); + let mut received = Vec::new(); + reader.read_to_end(&mut received).unwrap(); + received + }); + + let payload: Vec = (0..=u8::MAX).cycle().take(1024 * 1024).collect(); + let mut writer = RetryingWriter::new(PipeSink { pipe, waiting }); + writer.write_all(&payload).unwrap(); + writer.flush().unwrap(); + drop(writer); + + let received = reader_thread.join().unwrap(); + assert_eq!(received.len(), filled + payload.len()); + assert_eq!(&received[filled..], payload); + } +} From d3b072687a94969bb41589fa198966232bf938bf Mon Sep 17 00:00:00 2001 From: naokihaba <59875779+naokihaba@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:57:33 +0900 Subject: [PATCH 2/5] fix(dependencies): add "poll" feature to nix for Unix target --- crates/vt/Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/vt/Cargo.toml b/crates/vt/Cargo.toml index 7b6aed2b8..e85262e38 100644 --- a/crates/vt/Cargo.toml +++ b/crates/vt/Cargo.toml @@ -73,7 +73,7 @@ tempfile = { workspace = true } fspy = { workspace = true } [target.'cfg(unix)'.dependencies] -nix = { workspace = true, features = ["dir"] } +nix = { workspace = true, features = ["dir", "poll"] } [target.'cfg(windows)'.dependencies] winapi = { workspace = true, features = ["handleapi", "jobapi2", "winnt"] } From b0d48c776d7e49ab852dfe222a2db02b0166836e Mon Sep 17 00:00:00 2001 From: naokihaba <59875779+naokihaba@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:57:55 +0900 Subject: [PATCH 3/5] fix(reporter): use RetryingWriter for stdout and stderr in InterleavedLeafReporter --- crates/vt/src/session/reporter/interleaved/mod.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/crates/vt/src/session/reporter/interleaved/mod.rs b/crates/vt/src/session/reporter/interleaved/mod.rs index b91a1acda..7c23d98de 100644 --- a/crates/vt/src/session/reporter/interleaved/mod.rs +++ b/crates/vt/src/session/reporter/interleaved/mod.rs @@ -7,7 +7,7 @@ use vt_plan::{ExecutionItemDisplay, LeafExecutionKind}; use super::{ ColorSupport, ExitStatus, GraphExecutionReporter, GraphExecutionReporterBuilder, - LeafExecutionReporter, PipeWriters, StdioConfig, StdioSuggestion, + LeafExecutionReporter, PipeWriters, RetryingWriter, StdioConfig, StdioSuggestion, format_command_with_cache_status, maybe_strip_writer, write_leaf_trailing_output, }; use crate::session::event::{CacheStatus, CacheUpdateStatus, ExecutionError}; @@ -101,11 +101,11 @@ impl LeafExecutionReporter for InterleavedLeafReporter { suggestion: self.stdio_suggestion, writers: PipeWriters { stdout_writer: maybe_strip_writer( - Box::new(std::io::stdout()), + Box::new(RetryingWriter::new(std::io::stdout())), self.color_support.stdout, ), stderr_writer: maybe_strip_writer( - Box::new(std::io::stderr()), + Box::new(RetryingWriter::new(std::io::stderr())), self.color_support.stderr, ), }, From 09143761dd487c0c8f0febb092fd789f9f756568 Mon Sep 17 00:00:00 2001 From: naokihaba <59875779+naokihaba@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:58:05 +0900 Subject: [PATCH 4/5] fix(reporter): integrate RetryingWriter for stdout and stderr in LabeledLeafReporter --- crates/vt/src/session/reporter/labeled/mod.rs | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/crates/vt/src/session/reporter/labeled/mod.rs b/crates/vt/src/session/reporter/labeled/mod.rs index 3e91ac995..ae50f796e 100644 --- a/crates/vt/src/session/reporter/labeled/mod.rs +++ b/crates/vt/src/session/reporter/labeled/mod.rs @@ -7,7 +7,7 @@ use vt_plan::{ExecutionItemDisplay, LeafExecutionKind}; use super::{ ColorSupport, ExitStatus, GraphExecutionReporter, GraphExecutionReporterBuilder, - LeafExecutionReporter, PipeWriters, StdioConfig, StdioSuggestion, + LeafExecutionReporter, PipeWriters, RetryingWriter, StdioConfig, StdioSuggestion, format_command_with_cache_status, format_task_label, maybe_strip_writer, write_leaf_trailing_output, }; @@ -101,11 +101,17 @@ impl LeafExecutionReporter for LabeledLeafReporter { suggestion: StdioSuggestion::Piped, writers: PipeWriters { stdout_writer: Box::new(LabeledWriter::new( - maybe_strip_writer(Box::new(std::io::stdout()), self.color_support.stdout), + maybe_strip_writer( + Box::new(RetryingWriter::new(std::io::stdout())), + self.color_support.stdout, + ), prefix.as_bytes().to_vec(), )), stderr_writer: Box::new(LabeledWriter::new( - maybe_strip_writer(Box::new(std::io::stderr()), self.color_support.stderr), + maybe_strip_writer( + Box::new(RetryingWriter::new(std::io::stderr())), + self.color_support.stderr, + ), prefix.as_bytes().to_vec(), )), }, From 179f6bf1b2250e5a21587c2f813fc25148f7c3b4 Mon Sep 17 00:00:00 2001 From: naokihaba <59875779+naokihaba@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:58:13 +0900 Subject: [PATCH 5/5] fix(reporter): integrate RetryingWriter for stdout and stderr in PlainReporter --- crates/vt/src/session/reporter/mod.rs | 2 ++ crates/vt/src/session/reporter/plain.rs | 6 +++--- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/crates/vt/src/session/reporter/mod.rs b/crates/vt/src/session/reporter/mod.rs index b4c03fc99..af0396f93 100644 --- a/crates/vt/src/session/reporter/mod.rs +++ b/crates/vt/src/session/reporter/mod.rs @@ -27,6 +27,7 @@ mod grouped; mod interleaved; mod labeled; mod plain; +mod retrying_writer; pub mod summary; mod summary_reporter; @@ -37,6 +38,7 @@ pub use interleaved::InterleavedReporterBuilder; pub use labeled::LabeledReporterBuilder; use owo_colors::Style; pub use plain::PlainReporter; +use retrying_writer::RetryingWriter; pub use summary_reporter::SummaryReporterBuilder; use vt_path::AbsolutePath; use vt_plan::{ExecutionItemDisplay, LeafExecutionKind}; diff --git a/crates/vt/src/session/reporter/plain.rs b/crates/vt/src/session/reporter/plain.rs index 875db71b0..b23d53268 100644 --- a/crates/vt/src/session/reporter/plain.rs +++ b/crates/vt/src/session/reporter/plain.rs @@ -6,7 +6,7 @@ use std::io::Write; use super::{ - ColorSupport, LeafExecutionReporter, PipeWriters, StdioConfig, StdioSuggestion, + ColorSupport, LeafExecutionReporter, PipeWriters, RetryingWriter, StdioConfig, StdioSuggestion, format_cache_hit_message, format_error_message, maybe_strip_writer, }; // `maybe_strip_writer` is used for the child-process pipe writers; reporter @@ -85,11 +85,11 @@ impl LeafExecutionReporter for PlainReporter { suggestion: StdioSuggestion::Inherited, writers: PipeWriters { stdout_writer: maybe_strip_writer( - Box::new(std::io::stdout()), + Box::new(RetryingWriter::new(std::io::stdout())), self.color_support.stdout, ), stderr_writer: maybe_strip_writer( - Box::new(std::io::stderr()), + Box::new(RetryingWriter::new(std::io::stderr())), self.color_support.stderr, ), },