From 66f5bf1bbcfbc6ab65ee3ceb9d122a3e37641f79 Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Mon, 28 Apr 2025 21:55:31 +0800 Subject: [PATCH] tmp1 Signed-off-by: junxiang Mu <1948535941@qq.com> --- ecstore/src/erasure.rs | 21 ++++++++++++--------- ecstore/src/io.rs | 12 ++++++++---- ecstore/src/set_disk.rs | 17 +++++++---------- 3 files changed, 27 insertions(+), 23 deletions(-) diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index d5a08af4b..56f5959c5 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -1,17 +1,18 @@ use crate::bitrot::{BitrotReader, BitrotWriter}; use crate::error::clone_err; +use crate::io::Etag; use crate::quorum::{object_op_ignored_errs, reduce_write_quorum_errs}; use bytes::{Bytes, BytesMut}; use common::error::{Error, Result}; use futures::future::join_all; use reed_solomon_erasure::galois_8::ReedSolomon; use smallvec::SmallVec; -use tokio::sync::mpsc; use std::any::Any; use std::io::ErrorKind; use std::sync::Arc; use tokio::io::{AsyncRead, AsyncWrite}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::sync::mpsc; use tracing::warn; use tracing::{error, info}; // use tracing::debug; @@ -28,7 +29,7 @@ pub struct Erasure { encoder: Option, pub block_size: usize, _id: Uuid, - buf: Vec, + _buf: Vec, } impl Erasure { @@ -48,7 +49,7 @@ impl Erasure { block_size, encoder, _id: Uuid::new_v4(), - buf: vec![0u8; block_size], + _buf: vec![0u8; block_size], } } @@ -60,9 +61,9 @@ impl Erasure { // block_size: usize, total_size: usize, write_quorum: usize, - ) -> Result + ) -> Result<(usize, String)> where - S: AsyncRead + Unpin + Send + 'static, + S: AsyncRead + Etag + Unpin + Send + 'static, { // pin_mut!(body); // let mut reader = tokio_util::io::StreamReader::new( @@ -85,11 +86,11 @@ impl Erasure { remain } }; - + if new_len == 0 && total > 0 { break; } - + buf.resize(new_len, 0u8); match reader.read_exact(&mut buf).await { Ok(res) => res, @@ -106,9 +107,11 @@ impl Erasure { self_clone.clone().encode_data(&buf, &mut blocks)?; let _ = tx.send(blocks).await; } - Ok(total) + // let etag = reader.etag().await; + let etag = String::new(); + Ok((total, etag)) }); - + while let Some(blocks) = rx.recv().await { let write_futures = writers.iter_mut().enumerate().map(|(i, w_op)| { let i_inner = i; diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index 6bdbd629b..056260fdf 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -1,3 +1,4 @@ +use async_trait::async_trait; use bytes::Bytes; use futures::TryStreamExt; use md5::Digest; @@ -128,8 +129,9 @@ impl AsyncRead for HttpFileReader { } } -pub trait { - +#[async_trait] +pub trait Etag { + async fn etag(self) -> String; } pin_project! { @@ -140,7 +142,6 @@ pin_project! { } } - impl EtagReader { pub fn new(inner: R) -> Self { let (bytes_tx, mut bytes_rx) = mpsc::channel::(8); @@ -158,8 +159,11 @@ impl EtagReader { EtagReader { inner, bytes_tx, md5_rx } } +} - pub async fn etag(self) -> String { +#[async_trait] +impl Etag for EtagReader { + async fn etag(self) -> String { drop(self.inner); drop(self.bytes_tx); self.md5_rx.await.unwrap() diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 3c2b033c8..12491b105 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -3759,7 +3759,7 @@ impl ObjectIO for SetDisks { let tmp_object = format!("{}/{}/part.1", tmp_dir, fi.data_dir.unwrap()); - let mut erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); + let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); let is_inline_buffer = { if let Some(sc) = GLOBAL_StorageClass.get() { @@ -3798,11 +3798,11 @@ impl ObjectIO for SetDisks { } let stream = replace(&mut data.stream, Box::new(empty())); - let mut etag_stream = EtagReader::new(stream); + let etag_stream = EtagReader::new(stream); // TODO: etag from header - let w_size = Arc::new(erasure) + let (w_size, etag) = Arc::new(erasure) .encode(etag_stream, &mut writers, data.content_length, write_quorum) .await?; // TODO: 出错,删除临时目录 @@ -3810,7 +3810,6 @@ impl ObjectIO for SetDisks { error!("close_bitrot_writers err {:?}", err); } - let etag = etag_stream.etag().await; //TODO: userDefined user_defined.insert("etag".to_owned(), etag.clone()); @@ -4408,21 +4407,19 @@ impl StorageAPI for SetDisks { } } - let mut erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); + let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); let stream = replace(&mut data.stream, Box::new(empty())); - let mut etag_stream = EtagReader::new(stream); + let etag_stream = EtagReader::new(stream); - let w_size = erasure - .encode(&mut etag_stream, &mut writers, data.content_length, write_quorum) + let (w_size, mut etag) = Arc::new(erasure) + .encode(etag_stream, &mut writers, data.content_length, write_quorum) .await?; if let Err(err) = close_bitrot_writers(&mut writers).await { error!("close_bitrot_writers err {:?}", err); } - let mut etag = etag_stream.etag().await; - if let Some(ref tag) = opts.preserve_etag { etag = tag.clone(); }