diff --git a/ecstore/src/bitrot.rs b/ecstore/src/bitrot.rs index 1183f1ed5..8800c0643 100644 --- a/ecstore/src/bitrot.rs +++ b/ecstore/src/bitrot.rs @@ -1,13 +1,12 @@ use crate::{ - disk::{error::DiskError, DiskStore, FileReader, FileWriter}, + disk::{error::DiskError, DiskStore, FileReader, FileWriter, Reader}, erasure::{ReadAt, Write}, error::{Error, Result}, store_api::BitrotAlgorithm, - utils::crypto::hex, }; + use blake2::Blake2b512; use blake2::Digest as _; -use bytes::Bytes; use highway::{HighwayHash, HighwayHasher, Key}; use lazy_static::lazy_static; use sha2::{digest::core_api::BlockSizeUser, Digest, Sha256}; @@ -16,7 +15,7 @@ use std::{ collections::HashMap, io::{Cursor, Read}, }; -use tracing::warn; +use tracing::{error, warn}; use tokio::{ io::AsyncWriteExt, @@ -325,7 +324,8 @@ impl ReadAt for WholeBitrotReader { if self.buf.is_none() { let buf_len = self.till_offset - offset; let mut file = self.disk.read_file(&self.volume, &self.file_path).await?; - let (buf, _) = file.read_at(offset, buf_len).await?; + let mut buf = vec![0u8; buf_len]; + file.read_at(offset, &mut buf).await?; self.buf = Some(buf); } @@ -461,7 +461,8 @@ impl ReadAt for StreamingBitrotReader { let stream_offset = (offset / self.shard_size) * self.hasher.size() + offset; let buf_len = self.till_offset - stream_offset; let mut file = self.disk.read_file(&self.volume, &self.file_path).await?; - let (buf, _) = file.read_at(stream_offset, buf_len).await?; + let mut buf = vec![0u8; buf_len]; + file.read_at(stream_offset, &mut buf).await?; self.buf = buf; } if offset != self.curr_offset { @@ -538,6 +539,7 @@ struct BitrotFileReader { shard_size: usize, // buf: Vec, hash_bytes: Vec, + read_buf: Vec, } // fn ceil(a: usize, b: usize) -> usize { @@ -551,10 +553,11 @@ impl BitrotFileReader { inner, // till_offset: ceil(till_offset, shard_size) * hasher.size() + till_offset, curr_offset: 0, - hash_bytes: Vec::with_capacity(hasher.size()), + hash_bytes: vec![0u8; hasher.size()], hasher, shard_size, // buf: Vec::new(), + read_buf: Vec::new(), } } } @@ -563,15 +566,28 @@ impl BitrotFileReader { impl ReadAt for BitrotFileReader { async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec, usize)> { if offset % self.shard_size != 0 { + error!( + "BitrotFileReader read_at offset % self.shard_size != 0 , {} % {} = {}", + offset, + self.shard_size, + offset % self.shard_size + ); return Err(Error::new(DiskError::Unexpected)); } 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.read_buf.clear(); + self.read_buf.resize(buf_len, 0u8); + + self.inner.read_at(stream_offset, &mut self.read_buf).await?; + + let hash_bytes = &self.read_buf.as_slice()[0..self.hash_bytes.capacity()]; + + self.hash_bytes.clone_from_slice(hash_bytes); + let buf = self.read_buf.as_slice()[self.hash_bytes.capacity()..self.hash_bytes.capacity() + length].to_vec(); + self.hasher.reset(); self.hasher.update(&buf); let actual = self.hasher.clone().finalize(); @@ -584,32 +600,6 @@ 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 578a8476e..061f5a1a1 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -25,7 +25,16 @@ use protos::proto_gen::node_service::{ node_service_client::NodeServiceClient, ReadAtRequest, ReadAtResponse, WriteRequest, WriteResponse, }; use serde::{Deserialize, Serialize}; -use std::{any::Any, cmp::Ordering, collections::HashMap, fmt::Debug, io::SeekFrom, path::PathBuf, sync::Arc, usize}; +use std::{ + any::Any, + cmp::Ordering, + collections::HashMap, + fmt::Debug, + io::{Cursor, SeekFrom}, + path::PathBuf, + sync::Arc, + usize, +}; use time::OffsetDateTime; use tokio::{ fs::File, @@ -800,6 +809,13 @@ impl Write for RemoteFileWriter { } } +#[async_trait::async_trait] +pub trait Reader { + async fn read_at(&mut self, offset: usize, buf: &mut [u8]) -> Result; + async fn seek(&mut self, offset: usize) -> Result<()>; + async fn read_exact(&mut self, buf: &mut [u8]) -> Result; +} + #[derive(Debug)] pub enum FileReader { Local(LocalFileReader), @@ -808,60 +824,102 @@ pub enum FileReader { } #[async_trait::async_trait] -impl ReadAt for FileReader { - async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec, usize)> { +impl Reader for FileReader { + async fn read_at(&mut self, offset: usize, buf: &mut [u8]) -> Result { match self { - Self::Local(reader) => reader.read_at(offset, length).await, - Self::Remote(reader) => reader.read_at(offset, length).await, - Self::Buffer(reader) => reader.read_at(offset, length).await, + Self::Local(reader) => reader.read_at(offset, buf).await, + Self::Remote(reader) => reader.read_at(offset, buf).await, + Self::Buffer(reader) => reader.read_at(offset, buf).await, + } + } + async fn seek(&mut self, offset: usize) -> Result<()> { + match self { + Self::Local(reader) => reader.seek(offset).await, + Self::Remote(reader) => reader.seek(offset).await, + Self::Buffer(reader) => reader.seek(offset).await, + } + } + async fn read_exact(&mut self, buf: &mut [u8]) -> Result { + match self { + Self::Local(reader) => reader.read_exact(buf).await, + Self::Remote(reader) => reader.read_exact(buf).await, + Self::Buffer(reader) => reader.read_exact(buf).await, } } } #[derive(Debug)] pub struct BufferReader { - pub inner: Vec, + pub inner: Cursor>, + pos: usize, } impl BufferReader { pub fn new(inner: Vec) -> Self { - Self { inner } + Self { + inner: Cursor::new(inner), + pos: 0, + } } } #[async_trait::async_trait] -impl ReadAt for BufferReader { - async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec, usize)> { - let s = &self.inner[offset..offset + length]; - Ok((s.to_vec(), s.len())) +impl Reader for BufferReader { + #[tracing::instrument(level = "debug", skip(self, buf))] + async fn read_at(&mut self, offset: usize, buf: &mut [u8]) -> Result { + self.seek(offset).await?; + self.read_exact(buf).await + } + #[tracing::instrument(level = "debug", skip(self))] + async fn seek(&mut self, offset: usize) -> Result<()> { + if self.pos != offset { + self.inner.set_position(offset as u64); + } + + Ok(()) + } + #[tracing::instrument(level = "debug", skip(self))] + async fn read_exact(&mut self, buf: &mut [u8]) -> Result { + let bytes_read = self.inner.read_exact(buf).await?; + self.pos += buf.len(); + Ok(bytes_read) } } #[derive(Debug)] pub struct LocalFileReader { pub inner: File, + pos: usize, } impl LocalFileReader { pub fn new(inner: File) -> Self { - Self { inner } + Self { inner, pos: 0 } } } #[async_trait::async_trait] -impl ReadAt for LocalFileReader { - async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec, usize)> { - self.inner.seek(SeekFrom::Start(offset as u64)).await?; +impl Reader for LocalFileReader { + #[tracing::instrument(level = "debug", skip(self, buf))] + async fn read_at(&mut self, offset: usize, buf: &mut [u8]) -> Result { + self.seek(offset).await?; + self.read_exact(buf).await + } - let mut buffer = vec![0; length]; + #[tracing::instrument(level = "debug", skip(self))] + async fn seek(&mut self, offset: usize) -> Result<()> { + if self.pos != offset { + self.inner.seek(SeekFrom::Start(offset as u64)).await?; + self.pos = offset; + } - let bytes_read = self.inner.read_exact(&mut buffer).await?; - - // buffer.truncate(bytes_read); - - // warn!("LocalFileReader ReadAt need: {}, got: {}", length, bytes_read); - - Ok((buffer, bytes_read)) + Ok(()) + } + #[tracing::instrument(level = "debug", skip(self, buf))] + async fn read_exact(&mut self, buf: &mut [u8]) -> Result { + let bytes_read = self.inner.read_exact(buf).await?; + self.pos += buf.len(); + Ok(bytes_read) } } @@ -901,31 +959,38 @@ impl RemoteFileReader { } #[async_trait::async_trait] -impl ReadAt for RemoteFileReader { - async fn read_at(&mut self, offset: usize, length: usize) -> Result<(Vec, usize)> { - let request = ReadAtRequest { - disk: self.root.to_string_lossy().to_string(), - volume: self.volume.to_string(), - path: self.path.to_string(), - offset: offset.try_into().unwrap(), - length: length.try_into().unwrap(), - }; - self.tx.send(request).await?; +impl Reader for RemoteFileReader { + async fn read_at(&mut self, offset: usize, buf: &mut [u8]) -> Result { + unimplemented!() + // let request = ReadAtRequest { + // disk: self.root.to_string_lossy().to_string(), + // volume: self.volume.to_string(), + // path: self.path.to_string(), + // offset: offset.try_into().unwrap(), + // length: length.try_into().unwrap(), + // }; + // self.tx.send(request).await?; - if let Some(resp) = self.resp_stream.next().await { - let resp = resp?; - if resp.success { - info!("read at stream success"); - Ok((resp.data, resp.read_size.try_into().unwrap())) - } else { - let error_info = resp.error_info.unwrap_or("".to_string()); - info!("read at stream failed: {}", error_info); - Err(Error::from_string(error_info)) - } - } else { - let error_info = "can not get response"; - info!("read at stream failed: {}", error_info); - Err(Error::from_string(error_info)) - } + // if let Some(resp) = self.resp_stream.next().await { + // let resp = resp?; + // if resp.success { + // info!("read at stream success"); + // Ok((resp.data, resp.read_size.try_into().unwrap())) + // } else { + // let error_info = resp.error_info.unwrap_or("".to_string()); + // info!("read at stream failed: {}", error_info); + // Err(Error::from_string(error_info)) + // } + // } else { + // let error_info = "can not get response"; + // info!("read at stream failed: {}", error_info); + // Err(Error::from_string(error_info)) + // } + } + async fn seek(&mut self, offset: usize) -> Result<()> { + unimplemented!() + } + async fn read_exact(&mut self, buf: &mut [u8]) -> Result { + unimplemented!() } } diff --git a/rustfs/src/grpc.rs b/rustfs/src/grpc.rs index 3378c557e..bc8dd1d51 100644 --- a/rustfs/src/grpc.rs +++ b/rustfs/src/grpc.rs @@ -2,7 +2,7 @@ use std::{error::Error, io::ErrorKind, pin::Pin}; use ecstore::{ disk::{ - DeleteOptions, DiskInfoOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, UpdateMetadataOpts, + DeleteOptions, DiskInfoOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, Reader, UpdateMetadataOpts, WalkDirOptions, }, erasure::{ReadAt, Write}, @@ -624,13 +624,15 @@ impl Node for NodeService { } }; + let mut data = vec![0u8; v.length.try_into().unwrap()]; + match file_ref .as_mut() .unwrap() - .read_at(v.offset.try_into().unwrap(), v.length.try_into().unwrap()) + .read_at(v.offset.try_into().unwrap(), &mut data) .await { - Ok((data, read_size)) => tx.send(Ok(ReadAtResponse { + Ok(read_size) => tx.send(Ok(ReadAtResponse { success: true, data, read_size: read_size.try_into().unwrap(),