mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-26 05:56:50 +00:00
fix bitrot readat bug
This commit is contained in:
+46
-20
@@ -3,6 +3,7 @@ use crate::{
|
|||||||
erasure::{ReadAt, Write},
|
erasure::{ReadAt, Write},
|
||||||
error::{Error, Result},
|
error::{Error, Result},
|
||||||
store_api::BitrotAlgorithm,
|
store_api::BitrotAlgorithm,
|
||||||
|
utils::crypto::hex,
|
||||||
};
|
};
|
||||||
use blake2::Blake2b512;
|
use blake2::Blake2b512;
|
||||||
use blake2::Digest as _;
|
use blake2::Digest as _;
|
||||||
@@ -15,6 +16,7 @@ use std::{
|
|||||||
collections::HashMap,
|
collections::HashMap,
|
||||||
io::{Cursor, Read},
|
io::{Cursor, Read},
|
||||||
};
|
};
|
||||||
|
use tracing::warn;
|
||||||
|
|
||||||
use tokio::{
|
use tokio::{
|
||||||
io::AsyncWriteExt,
|
io::AsyncWriteExt,
|
||||||
@@ -516,7 +518,6 @@ impl Write for BitrotFileWriter {
|
|||||||
self.hasher.reset();
|
self.hasher.reset();
|
||||||
self.hasher.update(&buf);
|
self.hasher.update(&buf);
|
||||||
let hash_bytes = self.hasher.clone().finalize();
|
let hash_bytes = self.hasher.clone().finalize();
|
||||||
|
|
||||||
let _ = self.inner.write(&hash_bytes).await?;
|
let _ = self.inner.write(&hash_bytes).await?;
|
||||||
let _ = self.inner.write(buf).await?;
|
let _ = self.inner.write(buf).await?;
|
||||||
|
|
||||||
@@ -524,32 +525,36 @@ impl Write for BitrotFileWriter {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn new_bitrot_filewriter(inner: FileWriter, algo: BitrotAlgorithm, shard_size: usize) -> Result<BitrotWriter> {
|
pub fn new_bitrot_filewriter(inner: FileWriter, algo: BitrotAlgorithm, shard_size: usize) -> BitrotWriter {
|
||||||
Ok(Box::new(BitrotFileWriter::new(inner, algo, shard_size)))
|
Box::new(BitrotFileWriter::new(inner, algo, shard_size))
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
struct BitrotFileReader {
|
struct BitrotFileReader {
|
||||||
pub inner: FileReader,
|
pub inner: FileReader,
|
||||||
till_offset: usize,
|
// till_offset: usize,
|
||||||
curr_offset: usize,
|
curr_offset: usize,
|
||||||
hasher: Hasher,
|
hasher: Hasher,
|
||||||
shard_size: usize,
|
shard_size: usize,
|
||||||
buf: Vec<u8>,
|
// buf: Vec<u8>,
|
||||||
hash_bytes: Vec<u8>,
|
hash_bytes: Vec<u8>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// fn ceil(a: usize, b: usize) -> usize {
|
||||||
|
// (a + b - 1) / b
|
||||||
|
// }
|
||||||
|
|
||||||
impl BitrotFileReader {
|
impl BitrotFileReader {
|
||||||
pub fn new(inner: FileReader, algo: BitrotAlgorithm, till_offset: usize, shard_size: usize) -> Self {
|
pub fn new(inner: FileReader, algo: BitrotAlgorithm, _till_offset: usize, shard_size: usize) -> Self {
|
||||||
let hasher = algo.new();
|
let hasher = algo.new();
|
||||||
Self {
|
Self {
|
||||||
inner,
|
inner,
|
||||||
till_offset: till_offset.div_ceil(shard_size) * hasher.size() + till_offset,
|
// till_offset: ceil(till_offset, shard_size) * hasher.size() + till_offset,
|
||||||
curr_offset: 0,
|
curr_offset: 0,
|
||||||
hash_bytes: Vec::with_capacity(hasher.size()),
|
hash_bytes: Vec::with_capacity(hasher.size()),
|
||||||
hasher,
|
hasher,
|
||||||
shard_size,
|
shard_size,
|
||||||
buf: Vec::new(),
|
// buf: Vec::new(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -560,22 +565,17 @@ impl ReadAt for BitrotFileReader {
|
|||||||
if offset % self.shard_size != 0 {
|
if offset % self.shard_size != 0 {
|
||||||
return Err(Error::new(DiskError::Unexpected));
|
return Err(Error::new(DiskError::Unexpected));
|
||||||
}
|
}
|
||||||
if self.buf.is_empty() {
|
|
||||||
self.curr_offset = offset;
|
|
||||||
let stream_offset = (offset / self.shard_size) * self.hasher.size() + offset;
|
|
||||||
let buf_len = self.till_offset - stream_offset;
|
|
||||||
let (buf, _) = self.inner.read_at(stream_offset, buf_len).await?;
|
|
||||||
self.buf = buf;
|
|
||||||
}
|
|
||||||
if offset != self.curr_offset {
|
|
||||||
return Err(Error::new(DiskError::Unexpected));
|
|
||||||
}
|
|
||||||
|
|
||||||
self.hash_bytes = self.buf.drain(0..self.hash_bytes.capacity()).collect();
|
let stream_offset = (offset / self.shard_size) * self.hasher.size() + offset;
|
||||||
let buf = self.buf.drain(0..length).collect::<Vec<_>>();
|
let buf_len = self.hasher.size() + length;
|
||||||
|
let (mut read_buf, _) = self.inner.read_at(stream_offset, buf_len).await?;
|
||||||
|
|
||||||
|
self.hash_bytes = read_buf.drain(0..self.hash_bytes.capacity()).collect();
|
||||||
|
let buf = read_buf.drain(0..length).collect::<Vec<_>>();
|
||||||
self.hasher.reset();
|
self.hasher.reset();
|
||||||
self.hasher.update(&buf);
|
self.hasher.update(&buf);
|
||||||
let actual = self.hasher.clone().finalize();
|
let actual = self.hasher.clone().finalize();
|
||||||
|
|
||||||
if actual != self.hash_bytes {
|
if actual != self.hash_bytes {
|
||||||
return Err(Error::new(DiskError::FileCorrupt));
|
return Err(Error::new(DiskError::FileCorrupt));
|
||||||
}
|
}
|
||||||
@@ -584,6 +584,32 @@ impl ReadAt for BitrotFileReader {
|
|||||||
self.curr_offset += readed_len;
|
self.curr_offset += readed_len;
|
||||||
|
|
||||||
Ok((buf, readed_len))
|
Ok((buf, readed_len))
|
||||||
|
|
||||||
|
// if self.buf.is_empty() {
|
||||||
|
// self.curr_offset = offset;
|
||||||
|
// let stream_offset = (offset / self.shard_size) * self.hasher.size() + offset;
|
||||||
|
// let buf_len = self.till_offset - stream_offset;
|
||||||
|
// let (buf, _) = self.inner.read_at(stream_offset, buf_len).await?;
|
||||||
|
// self.buf = buf;
|
||||||
|
// }
|
||||||
|
// if offset != self.curr_offset {
|
||||||
|
// return Err(Error::new(DiskError::Unexpected));
|
||||||
|
// }
|
||||||
|
|
||||||
|
// self.hash_bytes = self.buf.drain(0..self.hash_bytes.capacity()).collect();
|
||||||
|
// let buf = self.buf.drain(0..length).collect::<Vec<_>>();
|
||||||
|
// self.hasher.reset();
|
||||||
|
// self.hasher.update(&buf);
|
||||||
|
// let actual = self.hasher.clone().finalize();
|
||||||
|
|
||||||
|
// if actual != self.hash_bytes {
|
||||||
|
// return Err(Error::new(DiskError::FileCorrupt));
|
||||||
|
// }
|
||||||
|
|
||||||
|
// let readed_len = buf.len();
|
||||||
|
// self.curr_offset += readed_len;
|
||||||
|
|
||||||
|
// Ok((buf, readed_len))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -34,8 +34,8 @@ use tokio::{
|
|||||||
};
|
};
|
||||||
use tokio_stream::wrappers::ReceiverStream;
|
use tokio_stream::wrappers::ReceiverStream;
|
||||||
use tonic::{service::interceptor::InterceptedService, transport::Channel, Request, Status, Streaming};
|
use tonic::{service::interceptor::InterceptedService, transport::Channel, Request, Status, Streaming};
|
||||||
use tracing::error;
|
|
||||||
use tracing::info;
|
use tracing::info;
|
||||||
|
use tracing::{error, warn};
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
pub type DiskStore = Arc<Box<dyn DiskAPI>>;
|
pub type DiskStore = Arc<Box<dyn DiskAPI>>;
|
||||||
@@ -855,9 +855,11 @@ impl ReadAt for LocalFileReader {
|
|||||||
|
|
||||||
let mut buffer = vec![0; length];
|
let mut buffer = vec![0; length];
|
||||||
|
|
||||||
let bytes_read = self.inner.read(&mut buffer).await?;
|
let bytes_read = self.inner.read_exact(&mut buffer).await?;
|
||||||
|
|
||||||
buffer.truncate(bytes_read);
|
// buffer.truncate(bytes_read);
|
||||||
|
|
||||||
|
// warn!("LocalFileReader ReadAt need: {}, got: {}", length, bytes_read);
|
||||||
|
|
||||||
Ok((buffer, bytes_read))
|
Ok((buffer, bytes_read))
|
||||||
}
|
}
|
||||||
|
|||||||
+14
-14
@@ -14,8 +14,8 @@ use tracing::warn;
|
|||||||
// use tracing::debug;
|
// use tracing::debug;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
use reader::reader::ChunkedStream;
|
// use reader::reader::ChunkedStream;
|
||||||
// use crate::chunk_stream::ChunkedStream;
|
use crate::chunk_stream::ChunkedStream;
|
||||||
use crate::disk::error::DiskError;
|
use crate::disk::error::DiskError;
|
||||||
|
|
||||||
pub struct Erasure {
|
pub struct Erasure {
|
||||||
@@ -58,8 +58,8 @@ impl Erasure {
|
|||||||
where
|
where
|
||||||
S: Stream<Item = Result<Bytes, StdError>> + Send + Sync,
|
S: Stream<Item = Result<Bytes, StdError>> + Send + Sync,
|
||||||
{
|
{
|
||||||
let stream = ChunkedStream::new(body, self.block_size);
|
// let stream = ChunkedStream::new(body, self.block_size);
|
||||||
// let mut stream = ChunkedStream::new(body, total_size, self.block_size, false);
|
let stream = ChunkedStream::new(body, total_size, self.block_size, false);
|
||||||
let mut total: usize = 0;
|
let mut total: usize = 0;
|
||||||
// let mut idx = 0;
|
// let mut idx = 0;
|
||||||
pin_mut!(stream);
|
pin_mut!(stream);
|
||||||
@@ -82,9 +82,8 @@ impl Erasure {
|
|||||||
|
|
||||||
let blocks = self.encode_data(data.as_ref())?;
|
let blocks = self.encode_data(data.as_ref())?;
|
||||||
|
|
||||||
// debug!(
|
// warn!(
|
||||||
// "encode shard {} size: {}/{} from block_size {}, total_size {} ",
|
// "encode shard size: {}/{} from block_size {}, total_size {} ",
|
||||||
// idx,
|
|
||||||
// blocks[0].len(),
|
// blocks[0].len(),
|
||||||
// blocks.len(),
|
// blocks.len(),
|
||||||
// data.len(),
|
// data.len(),
|
||||||
@@ -93,13 +92,14 @@ impl Erasure {
|
|||||||
|
|
||||||
let mut errs = Vec::new();
|
let mut errs = Vec::new();
|
||||||
|
|
||||||
for (i, w) in writers.iter_mut().enumerate() {
|
for (i, w_op) in writers.iter_mut().enumerate() {
|
||||||
if w.is_none() {
|
if let Some(w) = w_op {
|
||||||
continue;
|
match w.write(blocks[i].as_ref()).await {
|
||||||
}
|
Ok(_) => errs.push(None),
|
||||||
match w.as_mut().unwrap().write(blocks[i].as_ref()).await {
|
Err(e) => errs.push(Some(e)),
|
||||||
Ok(_) => errs.push(None),
|
}
|
||||||
Err(e) => errs.push(Some(e)),
|
} else {
|
||||||
|
errs.push(Some(Error::new(DiskError::DiskNotFound)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user