From ced17f6deb3d937ff16d315a9539d6b6a7114301 Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 3 Sep 2024 09:12:32 +0800 Subject: [PATCH] fix:reduce_write_quorum_errs --- ecstore/src/erasure.rs | 17 ++++++++--------- ecstore/src/quorum.rs | 6 +++++- 2 files changed, 13 insertions(+), 10 deletions(-) diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index 6c23842af..660849913 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -1,4 +1,5 @@ use crate::error::{Error, Result, StdError}; +use crate::quorum::{object_ignored_errs, reduce_write_quorum_errs}; use bytes::Bytes; use futures::future::join_all; use futures::{Stream, StreamExt}; @@ -49,7 +50,7 @@ impl Erasure { writers: &mut [W], // block_size: usize, total_size: usize, - _write_quorum: usize, + write_quorum: usize, ) -> Result where S: Stream> + Send + Sync + 'static, @@ -87,17 +88,15 @@ impl Erasure { for (i, w) in writers.iter_mut().enumerate() { match w.write_all(blocks[i].as_ref()).await { Ok(_) => errs.push(None), - Err(e) => errs.push(Some(e)), + Err(e) => errs.push(Some(Error::new(e))), } } - // debug!("{} encode_data write errs:{:?}", self.id, errs); - // // TODO: reduceWriteQuorumErrs - // for err in errs.iter() { - // if err.is_some() { - // return Err(Error::msg("message")); - // } - // } + let err_idx = reduce_write_quorum_errs(&errs, object_ignored_errs().as_slice(), write_quorum)?; + if errs[err_idx].is_some() { + let err = errs[err_idx].take().unwrap(); + return Err(err); + } } Err(e) => return Err(Error::from_std_error(e)), } diff --git a/ecstore/src/quorum.rs b/ecstore/src/quorum.rs index e50f0c82a..124ccfcba 100644 --- a/ecstore/src/quorum.rs +++ b/ecstore/src/quorum.rs @@ -15,6 +15,10 @@ pub fn is_file_not_found(e: &Error) -> bool { DiskError::FileNotFound.is(e) } +pub fn object_ignored_errs() -> Vec { + vec![is_file_not_found] +} + // 用于检查错误是否被忽略的函数 fn is_err_ignored(err: &Error, ignored_errs: &[CheckErrorFn]) -> bool { ignored_errs.iter().any(|&ignored_err| ignored_err(err)) @@ -76,7 +80,7 @@ fn reduce_quorum_errs(errs: &Vec>, ignored_errs: &[CheckErrorFn], // 根据读quorum验证错误数量 pub fn reduce_read_quorum_errs( - errs: &Vec>, + errs: &mut Vec>, ignored_errs: &[CheckErrorFn], read_quorum: usize, ) -> Result {