mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
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) <noreply@anthropic.com>
This commit is contained in:
1 parent
62481f1673
commit
f2498e122a
4 files changed
+88
-35
No files matched your search
@@ -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<usize> {
|
||||
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::*;
|
||||
|
||||
@@ -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<usize> {
|
||||
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).
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<us
|
||||
}
|
||||
}
|
||||
|
||||
/// Reads into `buf` at `offset` until it is full or the file ends, retrying
|
||||
/// interrupted reads; returns how many bytes were read.
|
||||
///
|
||||
/// Unlike [`read_at`], a count below `buf.len()` always means end of file.
|
||||
pub(crate) fn read_full_at(file: &File, buf: &mut [u8], offset: u64) -> io::Result<usize> {
|
||||
fill_at(|b, at| read_at(file, b, at), buf, offset)
|
||||
}
|
||||
|
||||
fn fill_at(
|
||||
mut read: impl FnMut(&mut [u8], u64) -> io::Result<usize>,
|
||||
buf: &mut [u8],
|
||||
offset: u64,
|
||||
) -> io::Result<usize> {
|
||||
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<usize> {
|
||||
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<u8> = (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<u8> = (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);
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user