From 6ab538050d288b03cb1a0fcbc0982a2a4b50c0b0 Mon Sep 17 00:00:00 2001 From: weisd Date: Sat, 2 Nov 2024 17:53:30 +0800 Subject: [PATCH] fix bitrot readat bug --- ecstore/src/bitrot.rs | 66 ++++++++++++++++++++++++++++------------- ecstore/src/disk/mod.rs | 8 +++-- ecstore/src/erasure.rs | 28 ++++++++--------- 3 files changed, 65 insertions(+), 37 deletions(-) diff --git a/ecstore/src/bitrot.rs b/ecstore/src/bitrot.rs index 03a3c6155..1183f1ed5 100644 --- a/ecstore/src/bitrot.rs +++ b/ecstore/src/bitrot.rs @@ -3,6 +3,7 @@ use crate::{ erasure::{ReadAt, Write}, error::{Error, Result}, store_api::BitrotAlgorithm, + utils::crypto::hex, }; use blake2::Blake2b512; use blake2::Digest as _; @@ -15,6 +16,7 @@ use std::{ collections::HashMap, io::{Cursor, Read}, }; +use tracing::warn; use tokio::{ io::AsyncWriteExt, @@ -516,7 +518,6 @@ impl Write for BitrotFileWriter { self.hasher.reset(); self.hasher.update(&buf); let hash_bytes = self.hasher.clone().finalize(); - let _ = self.inner.write(&hash_bytes).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 { - Ok(Box::new(BitrotFileWriter::new(inner, algo, shard_size))) +pub fn new_bitrot_filewriter(inner: FileWriter, algo: BitrotAlgorithm, shard_size: usize) -> BitrotWriter { + Box::new(BitrotFileWriter::new(inner, algo, shard_size)) } #[derive(Debug)] struct BitrotFileReader { pub inner: FileReader, - till_offset: usize, + // till_offset: usize, curr_offset: usize, hasher: Hasher, shard_size: usize, - buf: Vec, + // buf: Vec, hash_bytes: Vec, } +// fn ceil(a: usize, b: usize) -> usize { +// (a + b - 1) / b +// } + 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(); Self { 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, hash_bytes: Vec::with_capacity(hasher.size()), hasher, shard_size, - buf: Vec::new(), + // buf: Vec::new(), } } } @@ -560,22 +565,17 @@ impl ReadAt for BitrotFileReader { if offset % self.shard_size != 0 { 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 buf = self.buf.drain(0..length).collect::>(); + let stream_offset = (offset / self.shard_size) * self.hasher.size() + offset; + 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::>(); self.hasher.reset(); self.hasher.update(&buf); let actual = self.hasher.clone().finalize(); + if actual != self.hash_bytes { return Err(Error::new(DiskError::FileCorrupt)); } @@ -584,6 +584,32 @@ impl ReadAt for BitrotFileReader { self.curr_offset += 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::>(); + // 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)) } } diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 538e468fb..578a8476e 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -34,8 +34,8 @@ use tokio::{ }; use tokio_stream::wrappers::ReceiverStream; use tonic::{service::interceptor::InterceptedService, transport::Channel, Request, Status, Streaming}; -use tracing::error; use tracing::info; +use tracing::{error, warn}; use uuid::Uuid; pub type DiskStore = Arc>; @@ -855,9 +855,11 @@ impl ReadAt for LocalFileReader { 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)) } diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index 536d4acb3..4e27e0aec 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -14,8 +14,8 @@ use tracing::warn; // use tracing::debug; use uuid::Uuid; -use reader::reader::ChunkedStream; -// use crate::chunk_stream::ChunkedStream; +// use reader::reader::ChunkedStream; +use crate::chunk_stream::ChunkedStream; use crate::disk::error::DiskError; pub struct Erasure { @@ -58,8 +58,8 @@ impl Erasure { where S: Stream> + Send + Sync, { - 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, self.block_size); + let stream = ChunkedStream::new(body, total_size, self.block_size, false); let mut total: usize = 0; // let mut idx = 0; pin_mut!(stream); @@ -82,9 +82,8 @@ impl Erasure { let blocks = self.encode_data(data.as_ref())?; - // debug!( - // "encode shard {} size: {}/{} from block_size {}, total_size {} ", - // idx, + // warn!( + // "encode shard size: {}/{} from block_size {}, total_size {} ", // blocks[0].len(), // blocks.len(), // data.len(), @@ -93,13 +92,14 @@ impl Erasure { let mut errs = Vec::new(); - for (i, w) in writers.iter_mut().enumerate() { - if w.is_none() { - continue; - } - match w.as_mut().unwrap().write(blocks[i].as_ref()).await { - Ok(_) => errs.push(None), - Err(e) => errs.push(Some(e)), + for (i, w_op) in writers.iter_mut().enumerate() { + if let Some(w) = w_op { + match w.write(blocks[i].as_ref()).await { + Ok(_) => errs.push(None), + Err(e) => errs.push(Some(e)), + } + } else { + errs.push(Some(Error::new(DiskError::DiskNotFound))); } }