mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 00:38:16 +00:00
Optimization FileReader
This commit is contained in:
+26
-36
@@ -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<u8>,
|
||||
hash_bytes: Vec<u8>,
|
||||
read_buf: Vec<u8>,
|
||||
}
|
||||
|
||||
// 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<u8>, 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::<Vec<_>>();
|
||||
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::<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))
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+114
-49
@@ -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<usize>;
|
||||
async fn seek(&mut self, offset: usize) -> Result<()>;
|
||||
async fn read_exact(&mut self, buf: &mut [u8]) -> Result<usize>;
|
||||
}
|
||||
|
||||
#[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<u8>, usize)> {
|
||||
impl Reader for FileReader {
|
||||
async fn read_at(&mut self, offset: usize, buf: &mut [u8]) -> Result<usize> {
|
||||
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<usize> {
|
||||
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<u8>,
|
||||
pub inner: Cursor<Vec<u8>>,
|
||||
pos: usize,
|
||||
}
|
||||
|
||||
impl BufferReader {
|
||||
pub fn new(inner: Vec<u8>) -> 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<u8>, 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<usize> {
|
||||
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<usize> {
|
||||
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<u8>, 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<usize> {
|
||||
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<usize> {
|
||||
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<u8>, 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<usize> {
|
||||
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<usize> {
|
||||
unimplemented!()
|
||||
}
|
||||
}
|
||||
|
||||
+5
-3
@@ -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(),
|
||||
|
||||
Reference in New Issue
Block a user