From f2498e122afa516cac6cfc7584c29dcdc57275b8 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Thu, 1 Oct 2026 15:19:43 +0300 Subject: [PATCH] volume: read EC shards fully, like Go's ReadAt (#11537) EcVolumeShard::read_at and the scrub plan's EcLocalShard::read_at were a single pread/seek_read. That may legally return fewer bytes than asked mid-file (FUSE/NFS/CIFS mounts, a signal, very large requests), and an Interrupted error was not retried. Callers treat a short count as end of file or corruption: verify_ec_shards compared a zero tail and reported a parity mismatch, local scrub reported a broken shard, VolumeEcShardRead ended the stream early, and decode/rebuild/local needle reads failed. Add storage::io::read_full_at, which loops until the buffer is full or a read returns 0 and retries Interrupted, so a short count means EOF. Route both shard read_at methods through it, replace the encoder's private read_at_most with it, and reuse it for the Windows read_exact_at loop. Co-authored-by: Claude Opus 5.5 (1M context) --- .../src/storage/erasure_coding/ec_encoder.rs | 15 +-- .../src/storage/erasure_coding/ec_shard.rs | 4 +- .../src/storage/erasure_coding/ec_volume.rs | 2 +- seaweed-volume/src/storage/io.rs | 102 ++++++++++++++---- 4 files changed, 88 insertions(+), 35 deletions(-) diff --git a/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs b/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs index ded50f044..72d2cffa7 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs @@ -813,7 +813,7 @@ impl EncodeRun<'_> { // explicitly. for (i, buf) in self.buffers[..self.data_shards].iter_mut().enumerate() { let read_offset = offset + (i * block_size) as u64; - let n = read_at_most(self.dat_file, buf, read_offset)?; + let n = crate::storage::io::read_full_at(self.dat_file, buf, read_offset)?; buf[n..].fill(0); } @@ -834,19 +834,6 @@ impl EncodeRun<'_> { } } -/// Read into `buf` at `offset` until it is full or EOF; returns bytes read. -fn read_at_most(dat_file: &File, buf: &mut [u8], offset: u64) -> io::Result { - let mut n = 0; - while n < buf.len() { - let r = crate::storage::io::read_at(dat_file, &mut buf[n..], offset + n as u64)?; - if r == 0 { - break; - } - n += r; - } - Ok(n) -} - #[cfg(test)] mod tests { use super::*; diff --git a/seaweed-volume/src/storage/erasure_coding/ec_shard.rs b/seaweed-volume/src/storage/erasure_coding/ec_shard.rs index 36038f4a6..bda6e2b8a 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_shard.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_shard.rs @@ -87,14 +87,14 @@ impl EcVolumeShard { Ok(()) } - /// Read data at a specific offset. + /// Read data at a specific offset, filling `buf` unless the shard ends first. pub fn read_at(&self, buf: &mut [u8], offset: u64) -> io::Result { let file = self .ecd_file .as_ref() .ok_or_else(|| io::Error::other("shard file not open"))?; - crate::storage::io::read_at(file, buf, offset) + crate::storage::io::read_full_at(file, buf, offset) } /// Write data to the shard file (appends). diff --git a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs index eae9e7af4..de26e0ba9 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs @@ -4660,7 +4660,7 @@ impl EcLocalShard { .file .as_ref() .map_err(|e| io::Error::new(e.kind(), e.to_string()))?; - crate::storage::io::read_at(file, buf, offset) + crate::storage::io::read_full_at(file, buf, offset) } } diff --git a/seaweed-volume/src/storage/io.rs b/seaweed-volume/src/storage/io.rs index 93a29a326..57d32c44a 100644 --- a/seaweed-volume/src/storage/io.rs +++ b/seaweed-volume/src/storage/io.rs @@ -36,23 +36,11 @@ pub(crate) fn read_exact_at(file: &File, buf: &mut [u8], offset: u64) -> io::Res } #[cfg(windows)] { - use std::os::windows::fs::FileExt; - let mut filled = 0; - let mut at = offset; - while filled < buf.len() { - let n = match file.seek_read(&mut buf[filled..], at) { - Ok(n) => n, - Err(err) if err.kind() == io::ErrorKind::Interrupted => continue, - Err(err) => return Err(err), - }; - if n == 0 { - return Err(io::Error::new( - io::ErrorKind::UnexpectedEof, - "unexpected EOF in seek_read", - )); - } - filled += n; - at += n as u64; + if read_full_at(file, buf, offset)? < buf.len() { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "unexpected EOF in seek_read", + )); } } #[cfg(not(any(unix, windows)))] @@ -84,9 +72,34 @@ pub(crate) fn read_at(file: &File, buf: &mut [u8], offset: u64) -> io::Result io::Result { + fill_at(|b, at| read_at(file, b, at), buf, offset) +} + +fn fill_at( + mut read: impl FnMut(&mut [u8], u64) -> io::Result, + buf: &mut [u8], + offset: u64, +) -> io::Result { + let mut filled = 0; + while filled < buf.len() { + match read(&mut buf[filled..], offset + filled as u64) { + Ok(0) => break, + Ok(n) => filled += n, + Err(err) if err.kind() == io::ErrorKind::Interrupted => {} + Err(err) => return Err(err), + } + } + Ok(filled) +} + #[cfg(test)] mod tests { - use super::{read_at, read_exact_at}; + use super::{fill_at, read_at, read_exact_at, read_full_at}; use std::io::{ErrorKind, Write}; fn temp_file(bytes: &[u8]) -> tempfile::NamedTempFile { @@ -139,4 +152,57 @@ mod tests { let n = read_at(f.as_file(), &mut buf, 10).expect("read"); assert_eq!(n, 0); } + + /// A source that returns at most `chunk` bytes per call and fails with + /// `Interrupted` on its first call, like a network mount under a signal. + fn chunked(src: &[u8], chunk: usize) -> impl FnMut(&mut [u8], u64) -> std::io::Result { + let mut interrupted = false; + move |buf, at| { + if !interrupted { + interrupted = true; + return Err(ErrorKind::Interrupted.into()); + } + let at = (at as usize).min(src.len()); + let n = buf.len().min(chunk).min(src.len() - at); + buf[..n].copy_from_slice(&src[at..at + n]); + Ok(n) + } + } + + #[test] + fn fill_at_fills_across_short_and_interrupted_reads() { + let src: Vec = (0..=255).collect(); + let mut buf = [0u8; 100]; + let n = fill_at(chunked(&src, 7), &mut buf, 50).expect("read"); + assert_eq!(n, buf.len()); + assert_eq!(&buf[..], &src[50..150]); + } + + #[test] + fn fill_at_stops_at_end_of_source() { + let src: Vec = (0..=255).collect(); + let mut buf = [0u8; 100]; + let n = fill_at(chunked(&src, 7), &mut buf, 200).expect("read"); + assert_eq!(n, 56); + assert_eq!(&buf[..n], &src[200..]); + } + + #[test] + fn fill_at_propagates_other_errors() { + let mut buf = [0u8; 8]; + let err = fill_at(|_, _| Err(ErrorKind::PermissionDenied.into()), &mut buf, 0) + .expect_err("error"); + assert_eq!(err.kind(), ErrorKind::PermissionDenied); + } + + #[test] + fn read_full_at_returns_the_short_count_only_at_eof() { + let f = temp_file(b"0123456789"); + let mut buf = [0u8; 8]; + assert_eq!(read_full_at(f.as_file(), &mut buf, 0).expect("read"), 8); + assert_eq!(&buf, b"01234567"); + assert_eq!(read_full_at(f.as_file(), &mut buf, 6).expect("read"), 4); + assert_eq!(&buf[..4], b"6789"); + assert_eq!(read_full_at(f.as_file(), &mut buf, 10).expect("read"), 0); + } }