mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-05 12:57:42 +00:00
update ec share size
update bitrot
This commit is contained in:
+167
-841
File diff suppressed because it is too large
Load Diff
@@ -38,13 +38,13 @@ use crate::utils::path::{
|
||||
path_join, path_join_buf,
|
||||
};
|
||||
|
||||
use crate::erasure_coding::bitrot_verify;
|
||||
use common::defer;
|
||||
use path_absolutize::Absolutize;
|
||||
use rustfs_filemeta::{
|
||||
Cache, FileInfo, FileInfoOpts, FileMeta, MetaCacheEntry, MetacacheWriter, Opts, RawFileInfo, UpdateFn, get_file_info,
|
||||
read_xl_meta_no_data,
|
||||
};
|
||||
use rustfs_rio::bitrot_verify;
|
||||
use rustfs_utils::HashAlgorithm;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::fmt::Debug;
|
||||
|
||||
@@ -0,0 +1,458 @@
|
||||
use pin_project_lite::pin_project;
|
||||
use rustfs_utils::{HashAlgorithm, read_full, write_all};
|
||||
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite};
|
||||
|
||||
pin_project! {
|
||||
/// BitrotReader reads (hash+data) blocks from an async reader and verifies hash integrity.
|
||||
pub struct BitrotReader<R> {
|
||||
#[pin]
|
||||
inner: R,
|
||||
hash_algo: HashAlgorithm,
|
||||
shard_size: usize,
|
||||
buf: Vec<u8>,
|
||||
hash_buf: Vec<u8>,
|
||||
hash_read: usize,
|
||||
data_buf: Vec<u8>,
|
||||
data_read: usize,
|
||||
hash_checked: bool,
|
||||
}
|
||||
}
|
||||
|
||||
impl<R> BitrotReader<R>
|
||||
where
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
/// Create a new BitrotReader.
|
||||
pub fn new(inner: R, shard_size: usize, algo: HashAlgorithm) -> Self {
|
||||
let hash_size = algo.size();
|
||||
Self {
|
||||
inner,
|
||||
hash_algo: algo,
|
||||
shard_size,
|
||||
buf: Vec::new(),
|
||||
hash_buf: vec![0u8; hash_size],
|
||||
hash_read: 0,
|
||||
data_buf: Vec::new(),
|
||||
data_read: 0,
|
||||
hash_checked: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Read a single (hash+data) block, verify hash, and return the number of bytes read into `out`.
|
||||
/// Returns an error if hash verification fails or data exceeds shard_size.
|
||||
pub async fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> {
|
||||
if out.len() > self.shard_size {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidInput,
|
||||
format!("data size {} exceeds shard size {}", out.len(), self.shard_size),
|
||||
));
|
||||
}
|
||||
|
||||
let hash_size = self.hash_algo.size();
|
||||
// Read hash
|
||||
let mut hash_buf = vec![0u8; hash_size];
|
||||
if hash_size > 0 {
|
||||
self.inner.read_exact(&mut hash_buf).await?;
|
||||
}
|
||||
|
||||
let data_len = read_full(&mut self.inner, out).await?;
|
||||
|
||||
// // Read data
|
||||
// let mut data_len = 0;
|
||||
// while data_len < out.len() {
|
||||
// let n = self.inner.read(&mut out[data_len..]).await?;
|
||||
// if n == 0 {
|
||||
// break;
|
||||
// }
|
||||
// data_len += n;
|
||||
// // Only read up to one shard_size block
|
||||
// if data_len >= self.shard_size {
|
||||
// break;
|
||||
// }
|
||||
// }
|
||||
|
||||
if hash_size > 0 {
|
||||
let actual_hash = self.hash_algo.hash_encode(&out[..data_len]);
|
||||
if actual_hash != hash_buf {
|
||||
return Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "bitrot hash mismatch"));
|
||||
}
|
||||
}
|
||||
Ok(data_len)
|
||||
}
|
||||
}
|
||||
|
||||
pin_project! {
|
||||
/// BitrotWriter writes (hash+data) blocks to an async writer.
|
||||
pub struct BitrotWriter<W> {
|
||||
#[pin]
|
||||
inner: W,
|
||||
hash_algo: HashAlgorithm,
|
||||
shard_size: usize,
|
||||
buf: Vec<u8>,
|
||||
finished: bool,
|
||||
}
|
||||
}
|
||||
|
||||
impl<W> BitrotWriter<W>
|
||||
where
|
||||
W: AsyncWrite + Unpin + Send + Sync,
|
||||
{
|
||||
/// Create a new BitrotWriter.
|
||||
pub fn new(inner: W, shard_size: usize, algo: HashAlgorithm) -> Self {
|
||||
let hash_algo = algo;
|
||||
Self {
|
||||
inner,
|
||||
hash_algo,
|
||||
shard_size,
|
||||
buf: Vec::new(),
|
||||
finished: false,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn into_inner(self) -> W {
|
||||
self.inner
|
||||
}
|
||||
|
||||
/// Write a (hash+data) block. Returns the number of data bytes written.
|
||||
/// Returns an error if called after a short write or if data exceeds shard_size.
|
||||
pub async fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
if buf.is_empty() {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
if self.finished {
|
||||
return Err(std::io::Error::new(std::io::ErrorKind::InvalidInput, "bitrot writer already finished"));
|
||||
}
|
||||
|
||||
if buf.len() > self.shard_size {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidInput,
|
||||
format!("data size {} exceeds shard size {}", buf.len(), self.shard_size),
|
||||
));
|
||||
}
|
||||
|
||||
if buf.len() < self.shard_size {
|
||||
self.finished = true;
|
||||
}
|
||||
|
||||
let hash_algo = &self.hash_algo;
|
||||
|
||||
if hash_algo.size() > 0 {
|
||||
let hash = hash_algo.hash_encode(buf);
|
||||
self.buf.extend_from_slice(&hash);
|
||||
}
|
||||
|
||||
self.buf.extend_from_slice(buf);
|
||||
|
||||
// Write hash+data in one call
|
||||
let mut n = write_all(&mut self.inner, &self.buf).await?;
|
||||
|
||||
if n < hash_algo.size() {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::WriteZero,
|
||||
"short write: not enough bytes written",
|
||||
));
|
||||
}
|
||||
|
||||
n -= hash_algo.size();
|
||||
|
||||
self.buf.clear();
|
||||
|
||||
Ok(n)
|
||||
}
|
||||
}
|
||||
|
||||
pub fn bitrot_shard_file_size(size: usize, shard_size: usize, algo: HashAlgorithm) -> usize {
|
||||
if algo != HashAlgorithm::HighwayHash256S {
|
||||
return size;
|
||||
}
|
||||
size.div_ceil(shard_size) * algo.size() + size
|
||||
}
|
||||
|
||||
pub async fn bitrot_verify<R: AsyncRead + Unpin + Send>(
|
||||
mut r: R,
|
||||
want_size: usize,
|
||||
part_size: usize,
|
||||
algo: HashAlgorithm,
|
||||
_want: Vec<u8>,
|
||||
mut shard_size: usize,
|
||||
) -> std::io::Result<()> {
|
||||
let mut hash_buf = vec![0; algo.size()];
|
||||
let mut left = want_size;
|
||||
|
||||
if left != bitrot_shard_file_size(part_size, shard_size, algo.clone()) {
|
||||
return Err(std::io::Error::other("bitrot shard file size mismatch"));
|
||||
}
|
||||
|
||||
while left > 0 {
|
||||
let n = r.read_exact(&mut hash_buf).await?;
|
||||
left -= n;
|
||||
|
||||
if left < shard_size {
|
||||
shard_size = left;
|
||||
}
|
||||
|
||||
let mut buf = vec![0; shard_size];
|
||||
let read = r.read_exact(&mut buf).await?;
|
||||
|
||||
let actual_hash = algo.hash_encode(&buf);
|
||||
if actual_hash != hash_buf[0..n] {
|
||||
return Err(std::io::Error::other("bitrot hash mismatch"));
|
||||
}
|
||||
|
||||
left -= read;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Custom writer enum that supports inline buffer storage
|
||||
pub enum CustomWriter {
|
||||
/// Inline buffer writer - stores data in memory
|
||||
InlineBuffer(Vec<u8>),
|
||||
/// Disk-based writer using tokio file
|
||||
Other(Box<dyn AsyncWrite + Unpin + Send + Sync>),
|
||||
}
|
||||
|
||||
impl CustomWriter {
|
||||
/// Create a new inline buffer writer
|
||||
pub fn new_inline_buffer() -> Self {
|
||||
Self::InlineBuffer(Vec::new())
|
||||
}
|
||||
|
||||
/// Create a new disk writer from any AsyncWrite implementation
|
||||
pub fn new_tokio_writer<W>(writer: W) -> Self
|
||||
where
|
||||
W: AsyncWrite + Unpin + Send + Sync + 'static,
|
||||
{
|
||||
Self::Other(Box::new(writer))
|
||||
}
|
||||
|
||||
/// Get the inline buffer data if this is an inline buffer writer
|
||||
pub fn get_inline_data(&self) -> Option<&[u8]> {
|
||||
match self {
|
||||
Self::InlineBuffer(data) => Some(data),
|
||||
Self::Other(_) => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Extract the inline buffer data, consuming the writer
|
||||
pub fn into_inline_data(self) -> Option<Vec<u8>> {
|
||||
match self {
|
||||
Self::InlineBuffer(data) => Some(data),
|
||||
Self::Other(_) => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl AsyncWrite for CustomWriter {
|
||||
fn poll_write(
|
||||
self: std::pin::Pin<&mut Self>,
|
||||
cx: &mut std::task::Context<'_>,
|
||||
buf: &[u8],
|
||||
) -> std::task::Poll<std::io::Result<usize>> {
|
||||
match self.get_mut() {
|
||||
Self::InlineBuffer(data) => {
|
||||
data.extend_from_slice(buf);
|
||||
std::task::Poll::Ready(Ok(buf.len()))
|
||||
}
|
||||
Self::Other(writer) => {
|
||||
let pinned_writer = std::pin::Pin::new(writer.as_mut());
|
||||
pinned_writer.poll_write(cx, buf)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn poll_flush(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> std::task::Poll<std::io::Result<()>> {
|
||||
match self.get_mut() {
|
||||
Self::InlineBuffer(_) => std::task::Poll::Ready(Ok(())),
|
||||
Self::Other(writer) => {
|
||||
let pinned_writer = std::pin::Pin::new(writer.as_mut());
|
||||
pinned_writer.poll_flush(cx)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn poll_shutdown(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> std::task::Poll<std::io::Result<()>> {
|
||||
match self.get_mut() {
|
||||
Self::InlineBuffer(_) => std::task::Poll::Ready(Ok(())),
|
||||
Self::Other(writer) => {
|
||||
let pinned_writer = std::pin::Pin::new(writer.as_mut());
|
||||
pinned_writer.poll_shutdown(cx)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Wrapper around BitrotWriter that uses our custom writer
|
||||
pub struct BitrotWriterWrapper {
|
||||
bitrot_writer: BitrotWriter<CustomWriter>,
|
||||
writer_type: WriterType,
|
||||
}
|
||||
|
||||
/// Enum to track the type of writer we're using
|
||||
enum WriterType {
|
||||
InlineBuffer,
|
||||
Other,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for BitrotWriterWrapper {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("BitrotWriterWrapper")
|
||||
.field(
|
||||
"writer_type",
|
||||
&match self.writer_type {
|
||||
WriterType::InlineBuffer => "InlineBuffer",
|
||||
WriterType::Other => "Other",
|
||||
},
|
||||
)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
|
||||
impl BitrotWriterWrapper {
|
||||
/// Create a new BitrotWriterWrapper with custom writer
|
||||
pub fn new(writer: CustomWriter, shard_size: usize, checksum_algo: HashAlgorithm) -> Self {
|
||||
let writer_type = match &writer {
|
||||
CustomWriter::InlineBuffer(_) => WriterType::InlineBuffer,
|
||||
CustomWriter::Other(_) => WriterType::Other,
|
||||
};
|
||||
|
||||
Self {
|
||||
bitrot_writer: BitrotWriter::new(writer, shard_size, checksum_algo),
|
||||
writer_type,
|
||||
}
|
||||
}
|
||||
|
||||
/// Write data to the bitrot writer
|
||||
pub async fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
self.bitrot_writer.write(buf).await
|
||||
}
|
||||
|
||||
/// Extract the inline buffer data, consuming the wrapper
|
||||
pub fn into_inline_data(self) -> Option<Vec<u8>> {
|
||||
match self.writer_type {
|
||||
WriterType::InlineBuffer => {
|
||||
let writer = self.bitrot_writer.into_inner();
|
||||
writer.into_inline_data()
|
||||
}
|
||||
WriterType::Other => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
|
||||
use super::BitrotReader;
|
||||
use super::BitrotWriter;
|
||||
use rustfs_utils::HashAlgorithm;
|
||||
use std::io::Cursor;
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_bitrot_read_write_ok() {
|
||||
let data = b"hello world! this is a test shard.";
|
||||
let data_size = data.len();
|
||||
let shard_size = 8;
|
||||
|
||||
let buf: Vec<u8> = Vec::new();
|
||||
let writer = Cursor::new(buf);
|
||||
let mut bitrot_writer = BitrotWriter::new(writer, shard_size, HashAlgorithm::HighwayHash256);
|
||||
|
||||
let mut n = 0;
|
||||
for chunk in data.chunks(shard_size) {
|
||||
n += bitrot_writer.write(chunk).await.unwrap();
|
||||
}
|
||||
assert_eq!(n, data.len());
|
||||
|
||||
// 读
|
||||
let reader = bitrot_writer.into_inner();
|
||||
let mut bitrot_reader = BitrotReader::new(reader, shard_size, HashAlgorithm::HighwayHash256);
|
||||
let mut out = Vec::new();
|
||||
let mut n = 0;
|
||||
while n < data_size {
|
||||
let mut buf = vec![0u8; shard_size];
|
||||
let m = bitrot_reader.read(&mut buf).await.unwrap();
|
||||
assert_eq!(&buf[..m], &data[n..n + m]);
|
||||
|
||||
out.extend_from_slice(&buf[..m]);
|
||||
n += m;
|
||||
}
|
||||
|
||||
assert_eq!(n, data_size);
|
||||
assert_eq!(data, &out[..]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_bitrot_read_hash_mismatch() {
|
||||
let data = b"test data for bitrot";
|
||||
let data_size = data.len();
|
||||
let shard_size = 8;
|
||||
let buf: Vec<u8> = Vec::new();
|
||||
let writer = Cursor::new(buf);
|
||||
let mut bitrot_writer = BitrotWriter::new(writer, shard_size, HashAlgorithm::HighwayHash256);
|
||||
for chunk in data.chunks(shard_size) {
|
||||
let _ = bitrot_writer.write(chunk).await.unwrap();
|
||||
}
|
||||
let mut written = bitrot_writer.into_inner().into_inner();
|
||||
// change the last byte to make hash mismatch
|
||||
let pos = written.len() - 1;
|
||||
written[pos] ^= 0xFF;
|
||||
let reader = Cursor::new(written);
|
||||
let mut bitrot_reader = BitrotReader::new(reader, shard_size, HashAlgorithm::HighwayHash256);
|
||||
|
||||
let count = data_size.div_ceil(shard_size);
|
||||
|
||||
let mut idx = 0;
|
||||
let mut n = 0;
|
||||
while n < data_size {
|
||||
let mut buf = vec![0u8; shard_size];
|
||||
let res = bitrot_reader.read(&mut buf).await;
|
||||
|
||||
if idx == count - 1 {
|
||||
// 最后一个块,应该返回错误
|
||||
assert!(res.is_err());
|
||||
assert_eq!(res.unwrap_err().kind(), std::io::ErrorKind::InvalidData);
|
||||
break;
|
||||
}
|
||||
|
||||
let m = res.unwrap();
|
||||
|
||||
assert_eq!(&buf[..m], &data[n..n + m]);
|
||||
|
||||
n += m;
|
||||
idx += 1;
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_bitrot_read_write_none_hash() {
|
||||
let data = b"bitrot none hash test data!";
|
||||
let data_size = data.len();
|
||||
let shard_size = 8;
|
||||
|
||||
let buf: Vec<u8> = Vec::new();
|
||||
let writer = Cursor::new(buf);
|
||||
let mut bitrot_writer = BitrotWriter::new(writer, shard_size, HashAlgorithm::None);
|
||||
|
||||
let mut n = 0;
|
||||
for chunk in data.chunks(shard_size) {
|
||||
n += bitrot_writer.write(chunk).await.unwrap();
|
||||
}
|
||||
assert_eq!(n, data.len());
|
||||
|
||||
let reader = bitrot_writer.into_inner();
|
||||
let mut bitrot_reader = BitrotReader::new(reader, shard_size, HashAlgorithm::None);
|
||||
let mut out = Vec::new();
|
||||
let mut n = 0;
|
||||
while n < data_size {
|
||||
let mut buf = vec![0u8; shard_size];
|
||||
let m = bitrot_reader.read(&mut buf).await.unwrap();
|
||||
assert_eq!(&buf[..m], &data[n..n + m]);
|
||||
out.extend_from_slice(&buf[..m]);
|
||||
n += m;
|
||||
}
|
||||
assert_eq!(n, data_size);
|
||||
assert_eq!(data, &out[..]);
|
||||
}
|
||||
}
|
||||
@@ -1,18 +1,20 @@
|
||||
use super::BitrotReader;
|
||||
use super::Erasure;
|
||||
use crate::disk::error::Error;
|
||||
use crate::disk::error_reduce::reduce_errs;
|
||||
use futures::future::join_all;
|
||||
use pin_project_lite::pin_project;
|
||||
use rustfs_rio::BitrotReader;
|
||||
use std::io;
|
||||
use std::io::ErrorKind;
|
||||
use tokio::io::AsyncRead;
|
||||
use tokio::io::AsyncWrite;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use tracing::error;
|
||||
|
||||
pin_project! {
|
||||
pub(crate) struct ParallelReader {
|
||||
pub(crate) struct ParallelReader<R> {
|
||||
#[pin]
|
||||
readers: Vec<Option<BitrotReader>>,
|
||||
readers: Vec<Option<BitrotReader<R>>>,
|
||||
offset: usize,
|
||||
shard_size: usize,
|
||||
shard_file_size: usize,
|
||||
@@ -21,9 +23,12 @@ pub(crate) struct ParallelReader {
|
||||
}
|
||||
}
|
||||
|
||||
impl ParallelReader {
|
||||
impl<R> ParallelReader<R>
|
||||
where
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
// readers传入前应处理disk错误,确保每个reader达到可用数量的BitrotReader
|
||||
pub fn new(readers: Vec<Option<BitrotReader>>, e: Erasure, offset: usize, total_length: usize) -> Self {
|
||||
pub fn new(readers: Vec<Option<BitrotReader<R>>>, e: Erasure, offset: usize, total_length: usize) -> Self {
|
||||
let shard_size = e.shard_size();
|
||||
let shard_file_size = e.shard_file_size(total_length);
|
||||
|
||||
@@ -42,7 +47,10 @@ impl ParallelReader {
|
||||
}
|
||||
}
|
||||
|
||||
impl ParallelReader {
|
||||
impl<R> ParallelReader<R>
|
||||
where
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
pub async fn read(&mut self) -> (Vec<Option<Vec<u8>>>, Vec<Option<Error>>) {
|
||||
// if self.readers.len() != self.total_shards {
|
||||
// return Err(io::Error::new(ErrorKind::InvalidInput, "Invalid number of readers"));
|
||||
@@ -175,16 +183,17 @@ where
|
||||
}
|
||||
|
||||
impl Erasure {
|
||||
pub async fn decode<W>(
|
||||
pub async fn decode<W, R>(
|
||||
&self,
|
||||
writer: &mut W,
|
||||
readers: Vec<Option<BitrotReader>>,
|
||||
readers: Vec<Option<BitrotReader<R>>>,
|
||||
offset: usize,
|
||||
length: usize,
|
||||
total_length: usize,
|
||||
) -> (usize, Option<std::io::Error>)
|
||||
where
|
||||
W: tokio::io::AsyncWrite + Send + Sync + Unpin + 'static,
|
||||
W: AsyncWrite + Send + Sync + Unpin,
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
if readers.len() != self.data_shards + self.parity_shards {
|
||||
return (0, Some(io::Error::new(ErrorKind::InvalidInput, "Invalid number of readers")));
|
||||
|
||||
@@ -1,24 +1,22 @@
|
||||
use bytes::Bytes;
|
||||
use rustfs_rio::BitrotWriter;
|
||||
use rustfs_rio::Reader;
|
||||
// use std::io::Cursor;
|
||||
// use std::mem;
|
||||
use super::BitrotWriterWrapper;
|
||||
use super::Erasure;
|
||||
use crate::disk::error::Error;
|
||||
use crate::disk::error_reduce::count_errs;
|
||||
use crate::disk::error_reduce::{OBJECT_OP_IGNORED_ERRS, reduce_write_quorum_errs};
|
||||
use bytes::Bytes;
|
||||
use std::sync::Arc;
|
||||
use std::vec;
|
||||
use tokio::io::AsyncRead;
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
pub(crate) struct MultiWriter<'a> {
|
||||
writers: &'a mut [Option<BitrotWriter>],
|
||||
writers: &'a mut [Option<BitrotWriterWrapper>],
|
||||
write_quorum: usize,
|
||||
errs: Vec<Option<Error>>,
|
||||
}
|
||||
|
||||
impl<'a> MultiWriter<'a> {
|
||||
pub fn new(writers: &'a mut [Option<BitrotWriter>], write_quorum: usize) -> Self {
|
||||
pub fn new(writers: &'a mut [Option<BitrotWriterWrapper>], write_quorum: usize) -> Self {
|
||||
let length = writers.len();
|
||||
MultiWriter {
|
||||
writers,
|
||||
@@ -82,11 +80,11 @@ impl Erasure {
|
||||
pub async fn encode<R>(
|
||||
self: Arc<Self>,
|
||||
mut reader: R,
|
||||
writers: &mut [Option<BitrotWriter>],
|
||||
writers: &mut [Option<BitrotWriterWrapper>],
|
||||
quorum: usize,
|
||||
) -> std::io::Result<(R, usize)>
|
||||
where
|
||||
R: Reader + Send + Sync + Unpin + 'static,
|
||||
R: AsyncRead + Send + Sync + Unpin + 'static,
|
||||
{
|
||||
let (tx, mut rx) = mpsc::channel::<Vec<Bytes>>(8);
|
||||
|
||||
|
||||
@@ -11,16 +11,16 @@
|
||||
//! - **Compatibility**: Works with any shard size
|
||||
//! - **Use case**: Default behavior, recommended for most production use cases
|
||||
//!
|
||||
//! ### Hybrid Mode (`reed-solomon-simd` feature)
|
||||
//! - **Performance**: Uses SIMD optimization when possible, falls back to erasure implementation for small shards
|
||||
//! - **Compatibility**: Works with any shard size through intelligent fallback
|
||||
//! - **Reliability**: Best of both worlds - SIMD speed for large data, erasure stability for small data
|
||||
//! ### SIMD Mode (`reed-solomon-simd` feature)
|
||||
//! - **Performance**: Uses SIMD optimization for high-performance encoding/decoding
|
||||
//! - **Compatibility**: Works with any shard size through SIMD implementation
|
||||
//! - **Reliability**: High-performance SIMD implementation for large data processing
|
||||
//! - **Use case**: Use when maximum performance is needed for large data processing
|
||||
//!
|
||||
//! ## Feature Flags
|
||||
//!
|
||||
//! - Default: Use pure reed-solomon-erasure implementation (stable and reliable)
|
||||
//! - `reed-solomon-simd`: Use hybrid mode (SIMD + erasure fallback for optimal performance)
|
||||
//! - `reed-solomon-simd`: Use SIMD mode for optimal performance
|
||||
//! - `reed-solomon-erasure`: Explicitly enable pure erasure mode (same as default)
|
||||
//!
|
||||
//! ## Example
|
||||
@@ -38,25 +38,23 @@ use bytes::{Bytes, BytesMut};
|
||||
use reed_solomon_erasure::galois_8::ReedSolomon as ReedSolomonErasure;
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
use reed_solomon_simd;
|
||||
// use rustfs_rio::Reader;
|
||||
use smallvec::SmallVec;
|
||||
use std::io;
|
||||
use tokio::io::AsyncRead;
|
||||
use tracing::warn;
|
||||
use uuid::Uuid;
|
||||
|
||||
/// Reed-Solomon encoder variants supporting different implementations.
|
||||
#[allow(clippy::large_enum_variant)]
|
||||
pub enum ReedSolomonEncoder {
|
||||
/// Hybrid mode: SIMD with erasure fallback (when reed-solomon-simd feature is enabled)
|
||||
/// SIMD mode: High-performance SIMD implementation (when reed-solomon-simd feature is enabled)
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
Hybrid {
|
||||
SIMD {
|
||||
data_shards: usize,
|
||||
parity_shards: usize,
|
||||
// 使用RwLock确保线程安全,实现Send + Sync
|
||||
encoder_cache: std::sync::RwLock<Option<reed_solomon_simd::ReedSolomonEncoder>>,
|
||||
decoder_cache: std::sync::RwLock<Option<reed_solomon_simd::ReedSolomonDecoder>>,
|
||||
// erasure fallback for small shards or SIMD failures
|
||||
fallback_encoder: Box<ReedSolomonErasure>,
|
||||
},
|
||||
/// Pure erasure mode: default and when reed-solomon-erasure feature is specified
|
||||
Erasure(Box<ReedSolomonErasure>),
|
||||
@@ -66,18 +64,16 @@ impl Clone for ReedSolomonEncoder {
|
||||
fn clone(&self) -> Self {
|
||||
match self {
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
ReedSolomonEncoder::Hybrid {
|
||||
ReedSolomonEncoder::SIMD {
|
||||
data_shards,
|
||||
parity_shards,
|
||||
fallback_encoder,
|
||||
..
|
||||
} => ReedSolomonEncoder::Hybrid {
|
||||
} => ReedSolomonEncoder::SIMD {
|
||||
data_shards: *data_shards,
|
||||
parity_shards: *parity_shards,
|
||||
// 为新实例创建空的缓存,不共享缓存
|
||||
encoder_cache: std::sync::RwLock::new(None),
|
||||
decoder_cache: std::sync::RwLock::new(None),
|
||||
fallback_encoder: fallback_encoder.clone(),
|
||||
},
|
||||
ReedSolomonEncoder::Erasure(encoder) => ReedSolomonEncoder::Erasure(encoder.clone()),
|
||||
}
|
||||
@@ -89,18 +85,12 @@ impl ReedSolomonEncoder {
|
||||
pub fn new(data_shards: usize, parity_shards: usize) -> io::Result<Self> {
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
{
|
||||
// Hybrid mode: SIMD + erasure fallback when reed-solomon-simd feature is enabled
|
||||
let fallback_encoder = Box::new(
|
||||
ReedSolomonErasure::new(data_shards, parity_shards)
|
||||
.map_err(|e| io::Error::other(format!("Failed to create fallback erasure encoder: {:?}", e)))?,
|
||||
);
|
||||
|
||||
Ok(ReedSolomonEncoder::Hybrid {
|
||||
// SIMD mode when reed-solomon-simd feature is enabled
|
||||
Ok(ReedSolomonEncoder::SIMD {
|
||||
data_shards,
|
||||
parity_shards,
|
||||
encoder_cache: std::sync::RwLock::new(None),
|
||||
decoder_cache: std::sync::RwLock::new(None),
|
||||
fallback_encoder,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -117,11 +107,10 @@ impl ReedSolomonEncoder {
|
||||
pub fn encode(&self, shards: SmallVec<[&mut [u8]; 16]>) -> io::Result<()> {
|
||||
match self {
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
ReedSolomonEncoder::Hybrid {
|
||||
ReedSolomonEncoder::SIMD {
|
||||
data_shards,
|
||||
parity_shards,
|
||||
encoder_cache,
|
||||
fallback_encoder,
|
||||
..
|
||||
} => {
|
||||
let mut shards_vec: Vec<&mut [u8]> = shards.into_vec();
|
||||
@@ -129,30 +118,14 @@ impl ReedSolomonEncoder {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let shard_len = shards_vec[0].len();
|
||||
|
||||
// SIMD 性能最佳的最小 shard 大小 (通常 512-1024 字节)
|
||||
const SIMD_MIN_SHARD_SIZE: usize = 512;
|
||||
|
||||
// 如果 shard 太小,直接使用 fallback encoder
|
||||
if shard_len < SIMD_MIN_SHARD_SIZE {
|
||||
let fallback_shards: SmallVec<[&mut [u8]; 16]> = SmallVec::from_vec(shards_vec);
|
||||
return fallback_encoder
|
||||
.encode(fallback_shards)
|
||||
.map_err(|e| io::Error::other(format!("Fallback erasure encode error: {:?}", e)));
|
||||
}
|
||||
|
||||
// 尝试使用 SIMD,如果失败则回退到 fallback
|
||||
// 使用 SIMD 进行编码
|
||||
let simd_result = self.encode_with_simd(*data_shards, *parity_shards, encoder_cache, &mut shards_vec);
|
||||
|
||||
match simd_result {
|
||||
Ok(()) => Ok(()),
|
||||
Err(simd_error) => {
|
||||
warn!("SIMD encoding failed: {}, using fallback", simd_error);
|
||||
let fallback_shards: SmallVec<[&mut [u8]; 16]> = SmallVec::from_vec(shards_vec);
|
||||
fallback_encoder
|
||||
.encode(fallback_shards)
|
||||
.map_err(|e| io::Error::other(format!("Fallback erasure encode error: {:?}", e)))
|
||||
warn!("SIMD encoding failed: {}", simd_error);
|
||||
Err(simd_error)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -231,39 +204,20 @@ impl ReedSolomonEncoder {
|
||||
pub fn reconstruct(&self, shards: &mut [Option<Vec<u8>>]) -> io::Result<()> {
|
||||
match self {
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
ReedSolomonEncoder::Hybrid {
|
||||
ReedSolomonEncoder::SIMD {
|
||||
data_shards,
|
||||
parity_shards,
|
||||
decoder_cache,
|
||||
fallback_encoder,
|
||||
..
|
||||
} => {
|
||||
// Find a valid shard to determine length
|
||||
let shard_len = shards
|
||||
.iter()
|
||||
.find_map(|s| s.as_ref().map(|v| v.len()))
|
||||
.ok_or_else(|| io::Error::other("No valid shards found for reconstruction"))?;
|
||||
|
||||
// SIMD 性能最佳的最小 shard 大小
|
||||
const SIMD_MIN_SHARD_SIZE: usize = 512;
|
||||
|
||||
// 如果 shard 太小,直接使用 fallback encoder
|
||||
if shard_len < SIMD_MIN_SHARD_SIZE {
|
||||
return fallback_encoder
|
||||
.reconstruct(shards)
|
||||
.map_err(|e| io::Error::other(format!("Fallback erasure reconstruct error: {:?}", e)));
|
||||
}
|
||||
|
||||
// 尝试使用 SIMD,如果失败则回退到 fallback
|
||||
// 使用 SIMD 进行重构
|
||||
let simd_result = self.reconstruct_with_simd(*data_shards, *parity_shards, decoder_cache, shards);
|
||||
|
||||
match simd_result {
|
||||
Ok(()) => Ok(()),
|
||||
Err(simd_error) => {
|
||||
warn!("SIMD reconstruction failed: {}, using fallback", simd_error);
|
||||
fallback_encoder
|
||||
.reconstruct(shards)
|
||||
.map_err(|e| io::Error::other(format!("Fallback erasure reconstruct error: {:?}", e)))
|
||||
warn!("SIMD reconstruction failed: {}", simd_error);
|
||||
Err(simd_error)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -402,6 +356,10 @@ impl Clone for Erasure {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn calc_shard_size(block_size: usize, data_shards: usize) -> usize {
|
||||
(block_size.div_ceil(data_shards) + 1) & !1
|
||||
}
|
||||
|
||||
impl Erasure {
|
||||
/// Create a new Erasure instance.
|
||||
///
|
||||
@@ -439,7 +397,7 @@ impl Erasure {
|
||||
// let total_size = shard_size * self.total_shard_count();
|
||||
|
||||
// 数据切片数量
|
||||
let per_shard_size = data.len().div_ceil(self.data_shards);
|
||||
let per_shard_size = calc_shard_size(data.len(), self.data_shards);
|
||||
// 总需求大小
|
||||
let need_total_size = per_shard_size * self.total_shard_count();
|
||||
|
||||
@@ -507,7 +465,7 @@ impl Erasure {
|
||||
|
||||
/// Calculate the size of each shard.
|
||||
pub fn shard_size(&self) -> usize {
|
||||
self.block_size.div_ceil(self.data_shards)
|
||||
calc_shard_size(self.block_size, self.data_shards)
|
||||
}
|
||||
/// Calculate the total erasure file size for a given original size.
|
||||
// Returns the final erasure size from the original size
|
||||
@@ -518,7 +476,7 @@ impl Erasure {
|
||||
|
||||
let num_shards = total_length / self.block_size;
|
||||
let last_block_size = total_length % self.block_size;
|
||||
let last_shard_size = last_block_size.div_ceil(self.data_shards);
|
||||
let last_shard_size = calc_shard_size(last_block_size, self.data_shards);
|
||||
num_shards * self.shard_size() + last_shard_size
|
||||
}
|
||||
|
||||
@@ -536,22 +494,29 @@ impl Erasure {
|
||||
till_offset
|
||||
}
|
||||
|
||||
/// Encode all data from a rustfs_rio::Reader in blocks, calling an async callback for each encoded block.
|
||||
/// This method is async and returns the reader and total bytes read after all blocks are processed.
|
||||
/// Encode all data from a reader in blocks, calling an async callback for each encoded block.
|
||||
/// This method is async and returns the total bytes read after all blocks are processed.
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `reader` - A rustfs_rio::Reader to read data from.
|
||||
/// * `mut on_block` - Async callback: FnMut(Result<Vec<Bytes>, std::io::Error>) -> Future<Output=Result<(), E>> + Send
|
||||
/// * `reader` - An async reader implementing AsyncRead + Send + Sync + Unpin
|
||||
/// * `mut on_block` - Async callback that receives encoded blocks and returns a Result
|
||||
/// * `F` - Callback type: FnMut(Result<Vec<Bytes>, std::io::Error>) -> Future<Output=Result<(), E>> + Send
|
||||
/// * `Fut` - Future type returned by the callback
|
||||
/// * `E` - Error type returned by the callback
|
||||
/// * `R` - Reader type implementing AsyncRead + Send + Sync + Unpin
|
||||
///
|
||||
/// # Returns
|
||||
/// Result<(reader, total_bytes_read), E> after all data has been processed or on callback error.
|
||||
/// Result<usize, E> containing total bytes read, or error from callback
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns error if reading from reader fails or if callback returns error
|
||||
pub async fn encode_stream_callback_async<F, Fut, E, R>(
|
||||
self: std::sync::Arc<Self>,
|
||||
reader: &mut R,
|
||||
mut on_block: F,
|
||||
) -> Result<usize, E>
|
||||
where
|
||||
R: rustfs_rio::Reader + Send + Sync + Unpin,
|
||||
R: AsyncRead + Send + Sync + Unpin,
|
||||
F: FnMut(std::io::Result<Vec<Bytes>>) -> Fut + Send,
|
||||
Fut: std::future::Future<Output = Result<(), E>> + Send,
|
||||
{
|
||||
@@ -582,6 +547,7 @@ impl Erasure {
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
@@ -603,9 +569,9 @@ mod tests {
|
||||
// Case 5: total_length > block_size, aligned
|
||||
assert_eq!(erasure.shard_file_size(16), 4); // 16/8=2, last=0, 2*2+0=4
|
||||
|
||||
assert_eq!(erasure.shard_file_size(1248739), 312185); // 1248739/8=156092, last=3, 3 div_ceil 4=1, 156092*2+1=312185
|
||||
assert_eq!(erasure.shard_file_size(1248739), 312186); // 1248739/8=156092, last=3, 3 div_ceil 4=1, 156092*2+1=312185
|
||||
|
||||
assert_eq!(erasure.shard_file_size(43), 11); // 43/8=5, last=3, 3 div_ceil 4=1, 5*2+1=11
|
||||
assert_eq!(erasure.shard_file_size(43), 12); // 43/8=5, last=3, 3 div_ceil 4=1, 5*2+1=11
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -617,7 +583,7 @@ mod tests {
|
||||
#[cfg(not(feature = "reed-solomon-simd"))]
|
||||
let block_size = 8; // Pure erasure mode (default)
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
let block_size = 1024; // Hybrid mode - SIMD with fallback
|
||||
let block_size = 1024; // SIMD mode - SIMD with fallback
|
||||
|
||||
let erasure = Erasure::new(data_shards, parity_shards, block_size);
|
||||
|
||||
@@ -625,7 +591,7 @@ mod tests {
|
||||
#[cfg(not(feature = "reed-solomon-simd"))]
|
||||
let test_data = b"hello world".to_vec(); // Small data for erasure (default)
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
let test_data = b"Hybrid mode test data for encoding and decoding roundtrip verification with sufficient length to ensure shard size requirements are met for proper SIMD optimization.".repeat(20); // ~3KB for hybrid
|
||||
let test_data = b"SIMD mode test data for encoding and decoding roundtrip verification with sufficient length to ensure shard size requirements are met for proper SIMD optimization.".repeat(20); // ~3KB for SIMD
|
||||
|
||||
let data = &test_data;
|
||||
let encoded_shards = erasure.encode_data(data).unwrap();
|
||||
@@ -655,7 +621,7 @@ mod tests {
|
||||
|
||||
// Use different block sizes based on feature
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
let block_size = 512 * 3; // Hybrid mode - SIMD with fallback
|
||||
let block_size = 512 * 3; // SIMD mode
|
||||
#[cfg(not(feature = "reed-solomon-simd"))]
|
||||
let block_size = 8192; // Pure erasure mode (default)
|
||||
|
||||
@@ -700,7 +666,7 @@ mod tests {
|
||||
#[test]
|
||||
fn test_shard_size_and_file_size() {
|
||||
let erasure = Erasure::new(4, 2, 8);
|
||||
assert_eq!(erasure.shard_file_size(33), 9);
|
||||
assert_eq!(erasure.shard_file_size(33), 10);
|
||||
assert_eq!(erasure.shard_file_size(0), 0);
|
||||
}
|
||||
|
||||
@@ -722,7 +688,7 @@ mod tests {
|
||||
|
||||
// Use different block sizes based on feature
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
let block_size = 1024; // Hybrid mode
|
||||
let block_size = 1024; // SIMD mode
|
||||
#[cfg(not(feature = "reed-solomon-simd"))]
|
||||
let block_size = 8; // Pure erasure mode (default)
|
||||
|
||||
@@ -732,12 +698,12 @@ mod tests {
|
||||
let data =
|
||||
b"Async error test data with sufficient length to meet requirements for proper testing and validation.".repeat(20); // ~2KB
|
||||
|
||||
let mut rio_reader = Cursor::new(data);
|
||||
let mut reader = Cursor::new(data);
|
||||
let (tx, mut rx) = mpsc::channel::<Vec<Bytes>>(8);
|
||||
let erasure_clone = erasure.clone();
|
||||
let handle = tokio::spawn(async move {
|
||||
erasure_clone
|
||||
.encode_stream_callback_async::<_, _, (), _>(&mut rio_reader, move |res| {
|
||||
.encode_stream_callback_async::<_, _, (), _>(&mut reader, move |res| {
|
||||
let tx = tx.clone();
|
||||
async move {
|
||||
let shards = res.unwrap();
|
||||
@@ -765,7 +731,7 @@ mod tests {
|
||||
|
||||
// Use different block sizes based on feature
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
let block_size = 1024; // Hybrid mode
|
||||
let block_size = 1024; // SIMD mode
|
||||
#[cfg(not(feature = "reed-solomon-simd"))]
|
||||
let block_size = 8; // Pure erasure mode (default)
|
||||
|
||||
@@ -779,12 +745,12 @@ mod tests {
|
||||
// let data = b"callback".to_vec(); // 8 bytes to fit exactly in one 8-byte block
|
||||
|
||||
let data_clone = data.clone(); // Clone for later comparison
|
||||
let mut rio_reader = Cursor::new(data);
|
||||
let mut reader = Cursor::new(data);
|
||||
let (tx, mut rx) = mpsc::channel::<Vec<Bytes>>(8);
|
||||
let erasure_clone = erasure.clone();
|
||||
let handle = tokio::spawn(async move {
|
||||
erasure_clone
|
||||
.encode_stream_callback_async::<_, _, (), _>(&mut rio_reader, move |res| {
|
||||
.encode_stream_callback_async::<_, _, (), _>(&mut reader, move |res| {
|
||||
let tx = tx.clone();
|
||||
async move {
|
||||
let shards = res.unwrap();
|
||||
@@ -816,20 +782,20 @@ mod tests {
|
||||
assert_eq!(&recovered, &data_clone);
|
||||
}
|
||||
|
||||
// Tests specifically for hybrid mode (SIMD + erasure fallback)
|
||||
// Tests specifically for SIMD mode
|
||||
#[cfg(feature = "reed-solomon-simd")]
|
||||
mod hybrid_tests {
|
||||
mod simd_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_hybrid_encode_decode_roundtrip() {
|
||||
fn test_simd_encode_decode_roundtrip() {
|
||||
let data_shards = 4;
|
||||
let parity_shards = 2;
|
||||
let block_size = 1024; // Use larger block size for hybrid mode
|
||||
let block_size = 1024; // Use larger block size for SIMD mode
|
||||
let erasure = Erasure::new(data_shards, parity_shards, block_size);
|
||||
|
||||
// Use data that will create shards >= 512 bytes for SIMD optimization
|
||||
let test_data = b"Hybrid test data for encoding and decoding roundtrip verification with sufficient length to ensure shard size requirements are met for proper SIMD optimization and validation.";
|
||||
let test_data = b"SIMD mode test data for encoding and decoding roundtrip verification with sufficient length to ensure shard size requirements are met for proper SIMD optimization and validation.";
|
||||
let data = test_data.repeat(25); // Create much larger data: ~5KB total, ~1.25KB per shard
|
||||
|
||||
let encoded_shards = erasure.encode_data(&data).unwrap();
|
||||
@@ -854,10 +820,10 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_hybrid_all_zero_data() {
|
||||
fn test_simd_all_zero_data() {
|
||||
let data_shards = 4;
|
||||
let parity_shards = 2;
|
||||
let block_size = 1024; // Use larger block size for hybrid mode
|
||||
let block_size = 1024; // Use larger block size for SIMD mode
|
||||
let erasure = Erasure::new(data_shards, parity_shards, block_size);
|
||||
|
||||
// Create all-zero data that ensures adequate shard size for SIMD optimization
|
||||
@@ -997,39 +963,35 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_simd_smart_fallback() {
|
||||
fn test_simd_small_data_handling() {
|
||||
let data_shards = 4;
|
||||
let parity_shards = 2;
|
||||
let block_size = 32; // 很小的block_size,会导致小shard
|
||||
let block_size = 32; // Small block size for testing edge cases
|
||||
let erasure = Erasure::new(data_shards, parity_shards, block_size);
|
||||
|
||||
// 使用小数据,每个shard只有8字节,远小于512字节SIMD最小要求
|
||||
let small_data = b"tiny!123".to_vec(); // 8字节数据
|
||||
// Use small data to test SIMD handling of small shards
|
||||
let small_data = b"tiny!123".to_vec(); // 8 bytes data
|
||||
|
||||
// 应该能够成功编码(通过fallback)
|
||||
// Test encoding with small data
|
||||
let result = erasure.encode_data(&small_data);
|
||||
match result {
|
||||
Ok(shards) => {
|
||||
println!(
|
||||
"✅ Smart fallback worked: encoded {} bytes into {} shards",
|
||||
small_data.len(),
|
||||
shards.len()
|
||||
);
|
||||
println!("✅ SIMD encoding succeeded: {} bytes into {} shards", small_data.len(), shards.len());
|
||||
assert_eq!(shards.len(), data_shards + parity_shards);
|
||||
|
||||
// 测试解码
|
||||
// Test decoding
|
||||
let mut shards_opt: Vec<Option<Vec<u8>>> = shards.iter().map(|shard| Some(shard.to_vec())).collect();
|
||||
|
||||
// 丢失一些shard来测试恢复
|
||||
shards_opt[1] = None; // 丢失一个数据shard
|
||||
shards_opt[4] = None; // 丢失一个奇偶shard
|
||||
// Lose some shards to test recovery
|
||||
shards_opt[1] = None; // Lose one data shard
|
||||
shards_opt[4] = None; // Lose one parity shard
|
||||
|
||||
let decode_result = erasure.decode_data(&mut shards_opt);
|
||||
match decode_result {
|
||||
Ok(()) => {
|
||||
println!("✅ Smart fallback decode worked");
|
||||
println!("✅ SIMD decode worked");
|
||||
|
||||
// 验证恢复的数据
|
||||
// Verify recovered data
|
||||
let mut recovered = Vec::new();
|
||||
for shard in shards_opt.iter().take(data_shards) {
|
||||
recovered.extend_from_slice(shard.as_ref().unwrap());
|
||||
@@ -1038,17 +1000,17 @@ mod tests {
|
||||
println!("recovered: {:?}", recovered);
|
||||
println!("small_data: {:?}", small_data);
|
||||
assert_eq!(&recovered, &small_data);
|
||||
println!("✅ Data recovery successful with smart fallback");
|
||||
println!("✅ Data recovery successful with SIMD");
|
||||
}
|
||||
Err(e) => {
|
||||
println!("❌ Smart fallback decode failed: {}", e);
|
||||
// 对于很小的数据,如果decode失败也是可以接受的
|
||||
println!("❌ SIMD decode failed: {}", e);
|
||||
// For very small data, decode failure might be acceptable
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
println!("❌ Smart fallback encode failed: {}", e);
|
||||
// 如果连fallback都失败了,说明数据太小或配置有问题
|
||||
println!("❌ SIMD encode failed: {}", e);
|
||||
// For very small data or configuration issues, encoding might fail
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1143,14 +1105,14 @@ mod tests {
|
||||
let test_data = b"SIMD stream processing test with sufficient data length for multiple blocks and proper SIMD optimization verification!";
|
||||
let data = test_data.repeat(5); // Create owned Vec<u8>
|
||||
let data_clone = data.clone(); // Clone for later comparison
|
||||
let mut rio_reader = Cursor::new(data);
|
||||
let mut reader = Cursor::new(data);
|
||||
|
||||
let (tx, mut rx) = mpsc::channel::<Vec<Bytes>>(16);
|
||||
let erasure_clone = erasure.clone();
|
||||
|
||||
let handle = tokio::spawn(async move {
|
||||
erasure_clone
|
||||
.encode_stream_callback_async::<_, _, (), _>(&mut rio_reader, move |res| {
|
||||
.encode_stream_callback_async::<_, _, (), _>(&mut reader, move |res| {
|
||||
let tx = tx.clone();
|
||||
async move {
|
||||
let shards = res.unwrap();
|
||||
|
||||
@@ -1,19 +1,23 @@
|
||||
use super::BitrotReader;
|
||||
use super::BitrotWriterWrapper;
|
||||
use super::decode::ParallelReader;
|
||||
use crate::disk::error::{Error, Result};
|
||||
use crate::erasure_coding::encode::MultiWriter;
|
||||
use bytes::Bytes;
|
||||
use rustfs_rio::BitrotReader;
|
||||
use rustfs_rio::BitrotWriter;
|
||||
use tokio::io::AsyncRead;
|
||||
use tracing::info;
|
||||
|
||||
impl super::Erasure {
|
||||
pub async fn heal(
|
||||
pub async fn heal<R>(
|
||||
&self,
|
||||
writers: &mut [Option<BitrotWriter>],
|
||||
readers: Vec<Option<BitrotReader>>,
|
||||
writers: &mut [Option<BitrotWriterWrapper>],
|
||||
readers: Vec<Option<BitrotReader<R>>>,
|
||||
total_length: usize,
|
||||
_prefer: &[bool],
|
||||
) -> Result<()> {
|
||||
) -> Result<()>
|
||||
where
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
info!(
|
||||
"Erasure heal, writers len: {}, readers len: {}, total_length: {}",
|
||||
writers.len(),
|
||||
|
||||
@@ -3,4 +3,7 @@ pub mod encode;
|
||||
pub mod erasure;
|
||||
pub mod heal;
|
||||
|
||||
mod bitrot;
|
||||
pub use bitrot::*;
|
||||
|
||||
pub use erasure::{Erasure, ReedSolomonEncoder};
|
||||
|
||||
+1
-1
@@ -1,5 +1,5 @@
|
||||
pub mod admin_server_info;
|
||||
// pub mod bitrot;
|
||||
pub mod bitrot;
|
||||
pub mod bucket;
|
||||
pub mod cache_value;
|
||||
mod chunk_stream;
|
||||
|
||||
+168
-182
@@ -1,9 +1,11 @@
|
||||
use crate::bitrot::{create_bitrot_reader, create_bitrot_writer};
|
||||
use crate::disk::error_reduce::{OBJECT_OP_IGNORED_ERRS, reduce_read_quorum_errs, reduce_write_quorum_errs};
|
||||
use crate::disk::{
|
||||
self, CHECK_PART_DISK_NOT_FOUND, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS,
|
||||
conv_part_err_to_int, has_part_err,
|
||||
};
|
||||
use crate::erasure_coding;
|
||||
use crate::erasure_coding::bitrot_verify;
|
||||
use crate::error::{Error, Result};
|
||||
use crate::global::GLOBAL_MRFState;
|
||||
use crate::heal::data_usage_cache::DataUsageCache;
|
||||
@@ -64,7 +66,7 @@ use rustfs_filemeta::{
|
||||
FileInfo, FileMeta, FileMetaShallowVersion, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo,
|
||||
RawFileInfo, file_info_from_raw, merge_file_meta_versions,
|
||||
};
|
||||
use rustfs_rio::{BitrotReader, BitrotWriter, EtagResolvable, HashReader, Writer, bitrot_verify};
|
||||
use rustfs_rio::{EtagResolvable, HashReader};
|
||||
use rustfs_utils::HashAlgorithm;
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::hash::Hash;
|
||||
@@ -1865,51 +1867,31 @@ impl SetDisks {
|
||||
let mut readers = Vec::with_capacity(disks.len());
|
||||
let mut errors = Vec::with_capacity(disks.len());
|
||||
for (idx, disk_op) in disks.iter().enumerate() {
|
||||
if let Some(inline_data) = files[idx].data.clone() {
|
||||
let rd = Cursor::new(inline_data);
|
||||
let reader = BitrotReader::new(Box::new(rd), erasure.shard_size(), HashAlgorithm::HighwayHash256);
|
||||
readers.push(Some(reader));
|
||||
errors.push(None);
|
||||
} else if let Some(disk) = disk_op {
|
||||
// Calculate ceiling division of till_offset by shard_size
|
||||
let till_offset =
|
||||
till_offset.div_ceil(erasure.shard_size()) * HashAlgorithm::HighwayHash256.size() + till_offset;
|
||||
|
||||
let rd = disk
|
||||
.read_file_stream(
|
||||
bucket,
|
||||
&format!("{}/{}/part.{}", object, files[idx].data_dir.unwrap_or(Uuid::nil()), part_number),
|
||||
part_offset,
|
||||
till_offset,
|
||||
)
|
||||
.await?;
|
||||
let reader = BitrotReader::new(
|
||||
Box::new(rustfs_rio::WarpReader::new(rd)),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
);
|
||||
readers.push(Some(reader));
|
||||
errors.push(None);
|
||||
} else {
|
||||
errors.push(Some(DiskError::DiskNotFound));
|
||||
readers.push(None);
|
||||
match create_bitrot_reader(
|
||||
files[idx].data.as_deref(),
|
||||
disk_op.as_ref(),
|
||||
bucket,
|
||||
&format!("{}/{}/part.{}", object, files[idx].data_dir.unwrap_or_default(), part_number),
|
||||
part_offset,
|
||||
till_offset,
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(Some(reader)) => {
|
||||
readers.push(Some(reader));
|
||||
errors.push(None);
|
||||
}
|
||||
Ok(None) => {
|
||||
readers.push(None);
|
||||
errors.push(Some(DiskError::DiskNotFound));
|
||||
}
|
||||
Err(e) => {
|
||||
readers.push(None);
|
||||
errors.push(Some(e));
|
||||
}
|
||||
}
|
||||
|
||||
// if let Some(disk) = disk_op {
|
||||
// let checksum_info = files[idx].erasure.get_checksum_info(part_number);
|
||||
// let reader = new_bitrot_filereader(
|
||||
// disk.clone(),
|
||||
// files[idx].data.clone(),
|
||||
// bucket.to_owned(),
|
||||
// format!("{}/{}/part.{}", object, files[idx].data_dir.unwrap_or(Uuid::nil()), part_number),
|
||||
// till_offset,
|
||||
// checksum_info.algorithm,
|
||||
// erasure.shard_size(erasure.block_size),
|
||||
// );
|
||||
// readers.push(Some(reader));
|
||||
// } else {
|
||||
// readers.push(None)
|
||||
// }
|
||||
}
|
||||
|
||||
let nil_count = errors.iter().filter(|&e| e.is_none()).count();
|
||||
@@ -2478,59 +2460,31 @@ impl SetDisks {
|
||||
let mut prefer = vec![false; latest_disks.len()];
|
||||
for (index, disk) in latest_disks.iter().enumerate() {
|
||||
if let (Some(disk), Some(metadata)) = (disk, ©_parts_metadata[index]) {
|
||||
// let filereader = {
|
||||
// if let Some(ref data) = metadata.data {
|
||||
// Box::new(BufferReader::new(data.clone()))
|
||||
// } else {
|
||||
// let disk = disk.clone();
|
||||
// let part_path = format!("{}/{}/part.{}", object, src_data_dir, part.number);
|
||||
|
||||
// disk.read_file(bucket, &part_path).await?
|
||||
// }
|
||||
// };
|
||||
// let reader = new_bitrot_filereader(
|
||||
// disk.clone(),
|
||||
// metadata.data.clone(),
|
||||
// bucket.to_owned(),
|
||||
// format!("{}/{}/part.{}", object, src_data_dir, part.number),
|
||||
// till_offset,
|
||||
// checksum_algo.clone(),
|
||||
// erasure.shard_size(erasure.block_size),
|
||||
// );
|
||||
|
||||
if let Some(ref data) = metadata.data {
|
||||
let rd = Cursor::new(data.clone());
|
||||
let reader =
|
||||
BitrotReader::new(Box::new(rd), erasure.shard_size(), checksum_algo.clone());
|
||||
readers.push(Some(reader));
|
||||
// errors.push(None);
|
||||
} else {
|
||||
let length =
|
||||
till_offset.div_ceil(erasure.shard_size()) * checksum_algo.size() + till_offset;
|
||||
let rd = match disk
|
||||
.read_file_stream(
|
||||
bucket,
|
||||
&format!("{}/{}/part.{}", object, src_data_dir, part.number),
|
||||
0,
|
||||
length,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(rd) => rd,
|
||||
Err(e) => {
|
||||
// errors.push(Some(e.into()));
|
||||
error!("heal_object read_file err: {:?}", e);
|
||||
writers.push(None);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let reader = BitrotReader::new(
|
||||
Box::new(rustfs_rio::WarpReader::new(rd)),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
);
|
||||
readers.push(Some(reader));
|
||||
// errors.push(None);
|
||||
match create_bitrot_reader(
|
||||
metadata.data.as_deref(),
|
||||
Some(disk),
|
||||
bucket,
|
||||
&format!("{}/{}/part.{}", object, src_data_dir, part.number),
|
||||
0,
|
||||
till_offset,
|
||||
erasure.shard_size(),
|
||||
checksum_algo.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(Some(reader)) => {
|
||||
readers.push(Some(reader));
|
||||
}
|
||||
Ok(None) => {
|
||||
error!("heal_object disk not available");
|
||||
readers.push(None);
|
||||
continue;
|
||||
}
|
||||
Err(e) => {
|
||||
error!("heal_object read_file err: {:?}", e);
|
||||
readers.push(None);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
prefer[index] = disk.host_name().is_empty();
|
||||
@@ -2549,55 +2503,67 @@ impl SetDisks {
|
||||
};
|
||||
|
||||
for disk in out_dated_disks.iter() {
|
||||
if let Some(disk) = disk {
|
||||
// let filewriter = {
|
||||
// if is_inline_buffer {
|
||||
// Box::new(Cursor::new(Vec::new()))
|
||||
// } else {
|
||||
// let disk = disk.clone();
|
||||
// let part_path = format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number);
|
||||
// disk.create_file("", RUSTFS_META_TMP_BUCKET, &part_path, 0).await?
|
||||
// }
|
||||
// };
|
||||
let writer = create_bitrot_writer(
|
||||
is_inline_buffer,
|
||||
disk.as_ref(),
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
&format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number),
|
||||
erasure.shard_file_size(part.size),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
)
|
||||
.await?;
|
||||
writers.push(Some(writer));
|
||||
|
||||
if is_inline_buffer {
|
||||
let writer = BitrotWriter::new(
|
||||
Writer::from_cursor(Cursor::new(Vec::new())),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
);
|
||||
writers.push(Some(writer));
|
||||
} else {
|
||||
let f = disk
|
||||
.create_file(
|
||||
"",
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
&format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number),
|
||||
0,
|
||||
)
|
||||
.await?;
|
||||
let writer = BitrotWriter::new(
|
||||
Writer::from_tokio_writer(f),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
);
|
||||
writers.push(Some(writer));
|
||||
}
|
||||
// if let Some(disk) = disk {
|
||||
// // let filewriter = {
|
||||
// // if is_inline_buffer {
|
||||
// // Box::new(Cursor::new(Vec::new()))
|
||||
// // } else {
|
||||
// // let disk = disk.clone();
|
||||
// // let part_path = format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number);
|
||||
// // disk.create_file("", RUSTFS_META_TMP_BUCKET, &part_path, 0).await?
|
||||
// // }
|
||||
// // };
|
||||
|
||||
// let writer = new_bitrot_filewriter(
|
||||
// disk.clone(),
|
||||
// RUSTFS_META_TMP_BUCKET,
|
||||
// format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number).as_str(),
|
||||
// is_inline_buffer,
|
||||
// DEFAULT_BITROT_ALGO,
|
||||
// erasure.shard_size(erasure.block_size),
|
||||
// )
|
||||
// .await?;
|
||||
// if is_inline_buffer {
|
||||
// let writer = BitrotWriter::new(
|
||||
// Writer::from_cursor(Cursor::new(Vec::new())),
|
||||
// erasure.shard_size(),
|
||||
// HashAlgorithm::HighwayHash256,
|
||||
// );
|
||||
// writers.push(Some(writer));
|
||||
// } else {
|
||||
// let f = disk
|
||||
// .create_file(
|
||||
// "",
|
||||
// RUSTFS_META_TMP_BUCKET,
|
||||
// &format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number),
|
||||
// 0,
|
||||
// )
|
||||
// .await?;
|
||||
// let writer = BitrotWriter::new(
|
||||
// Writer::from_tokio_writer(f),
|
||||
// erasure.shard_size(),
|
||||
// HashAlgorithm::HighwayHash256,
|
||||
// );
|
||||
// writers.push(Some(writer));
|
||||
// }
|
||||
|
||||
// writers.push(Some(writer));
|
||||
} else {
|
||||
writers.push(None);
|
||||
}
|
||||
// // let writer = new_bitrot_filewriter(
|
||||
// // disk.clone(),
|
||||
// // RUSTFS_META_TMP_BUCKET,
|
||||
// // format!("{}/{}/part.{}", tmp_id, dst_data_dir, part.number).as_str(),
|
||||
// // is_inline_buffer,
|
||||
// // DEFAULT_BITROT_ALGO,
|
||||
// // erasure.shard_size(erasure.block_size),
|
||||
// // )
|
||||
// // .await?;
|
||||
|
||||
// // writers.push(Some(writer));
|
||||
// } else {
|
||||
// writers.push(None);
|
||||
// }
|
||||
}
|
||||
|
||||
// Heal each part. erasure.Heal() will write the healed
|
||||
@@ -2630,8 +2596,7 @@ impl SetDisks {
|
||||
// if let Some(w) = writer.as_any().downcast_ref::<BitrotFileWriter>() {
|
||||
// parts_metadata[index].data = Some(w.inline_data().to_vec());
|
||||
// }
|
||||
parts_metadata[index].data =
|
||||
Some(writer.into_inner().into_cursor_inner().unwrap_or_default());
|
||||
parts_metadata[index].data = Some(writer.into_inline_data().unwrap_or_default());
|
||||
}
|
||||
parts_metadata[index].set_inline_data();
|
||||
} else {
|
||||
@@ -3882,27 +3847,38 @@ impl ObjectIO for SetDisks {
|
||||
let mut errors = Vec::with_capacity(shuffle_disks.len());
|
||||
for disk_op in shuffle_disks.iter() {
|
||||
if let Some(disk) = disk_op {
|
||||
let writer = if is_inline_buffer {
|
||||
BitrotWriter::new(
|
||||
Writer::from_cursor(Cursor::new(Vec::new())),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
)
|
||||
} else {
|
||||
let f = match disk
|
||||
.create_file("", RUSTFS_META_TMP_BUCKET, &tmp_object, erasure.shard_file_size(data.content_length))
|
||||
.await
|
||||
{
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
errors.push(Some(e));
|
||||
writers.push(None);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let writer = create_bitrot_writer(
|
||||
is_inline_buffer,
|
||||
Some(disk),
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
&tmp_object,
|
||||
erasure.shard_file_size(data.content_length),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
)
|
||||
.await?;
|
||||
|
||||
BitrotWriter::new(Writer::from_tokio_writer(f), erasure.shard_size(), HashAlgorithm::HighwayHash256)
|
||||
};
|
||||
// let writer = if is_inline_buffer {
|
||||
// BitrotWriter::new(
|
||||
// Writer::from_cursor(Cursor::new(Vec::new())),
|
||||
// erasure.shard_size(),
|
||||
// HashAlgorithm::HighwayHash256,
|
||||
// )
|
||||
// } else {
|
||||
// let f = match disk
|
||||
// .create_file("", RUSTFS_META_TMP_BUCKET, &tmp_object, erasure.shard_file_size(data.content_length))
|
||||
// .await
|
||||
// {
|
||||
// Ok(f) => f,
|
||||
// Err(e) => {
|
||||
// errors.push(Some(e));
|
||||
// writers.push(None);
|
||||
// continue;
|
||||
// }
|
||||
// };
|
||||
|
||||
// BitrotWriter::new(Writer::from_tokio_writer(f), erasure.shard_size(), HashAlgorithm::HighwayHash256)
|
||||
// };
|
||||
|
||||
writers.push(Some(writer));
|
||||
errors.push(None);
|
||||
@@ -3952,7 +3928,7 @@ impl ObjectIO for SetDisks {
|
||||
for (i, fi) in parts_metadatas.iter_mut().enumerate() {
|
||||
if is_inline_buffer {
|
||||
if let Some(writer) = writers[i].take() {
|
||||
fi.data = Some(writer.into_inner().into_cursor_inner().unwrap_or_default());
|
||||
fi.data = Some(writer.into_inline_data().unwrap_or_default());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4553,21 +4529,32 @@ impl StorageAPI for SetDisks {
|
||||
let mut errors = Vec::with_capacity(shuffle_disks.len());
|
||||
for disk_op in shuffle_disks.iter() {
|
||||
if let Some(disk) = disk_op {
|
||||
let writer = {
|
||||
let f = match disk
|
||||
.create_file("", RUSTFS_META_TMP_BUCKET, &tmp_part_path, erasure.shard_file_size(data.content_length))
|
||||
.await
|
||||
{
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
errors.push(Some(e));
|
||||
writers.push(None);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let writer = create_bitrot_writer(
|
||||
false,
|
||||
Some(disk),
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
&tmp_part_path,
|
||||
erasure.shard_file_size(data.content_length),
|
||||
erasure.shard_size(),
|
||||
HashAlgorithm::HighwayHash256,
|
||||
)
|
||||
.await?;
|
||||
|
||||
BitrotWriter::new(Writer::from_tokio_writer(f), erasure.shard_size(), HashAlgorithm::HighwayHash256)
|
||||
};
|
||||
// let writer = {
|
||||
// let f = match disk
|
||||
// .create_file("", RUSTFS_META_TMP_BUCKET, &tmp_part_path, erasure.shard_file_size(data.content_length))
|
||||
// .await
|
||||
// {
|
||||
// Ok(f) => f,
|
||||
// Err(e) => {
|
||||
// errors.push(Some(e));
|
||||
// writers.push(None);
|
||||
// continue;
|
||||
// }
|
||||
// };
|
||||
|
||||
// BitrotWriter::new(Writer::from_tokio_writer(f), erasure.shard_size(), HashAlgorithm::HighwayHash256)
|
||||
// };
|
||||
|
||||
writers.push(Some(writer));
|
||||
errors.push(None);
|
||||
@@ -6079,7 +6066,6 @@ mod tests {
|
||||
// Test object directory dangling detection
|
||||
let errs = vec![Some(DiskError::FileNotFound), Some(DiskError::FileNotFound), None];
|
||||
assert!(is_object_dir_dang_ling(&errs));
|
||||
|
||||
let errs2 = vec![None, None, None];
|
||||
assert!(!is_object_dir_dang_ling(&errs2));
|
||||
|
||||
|
||||
Reference in New Issue
Block a user