diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index 3ca27fe0d..7ecf8f84c 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -1,6 +1,8 @@ +use bytes::Bytes; use futures::TryStreamExt; use md5::Digest; use md5::Md5; +use tokio::sync::mpsc; use std::pin::Pin; use std::task::Context; use std::task::Poll; @@ -125,33 +127,51 @@ impl AsyncRead for HttpFileReader { pub struct EtagReader { inner: R, - md5: Md5, + bytes_tx: mpsc::Sender, + md5_rx: oneshot::Receiver, } impl EtagReader { pub fn new(inner: R) -> Self { - EtagReader { inner, md5: Md5::new() } + let (bytes_tx, mut bytes_rx) = mpsc::channel::(8); + let (md5_tx, md5_rx) = oneshot::channel::(); + + tokio::task::spawn_blocking(move || { + let mut md5 = Md5::new(); + while let Some(bytes) = bytes_rx.blocking_recv() { + md5.update(&bytes); + } + let digest = md5.finalize(); + let etag = hex_simd::encode_to_string(digest, hex_simd::AsciiCase::Lower); + let _ = md5_tx.send(etag); + }); + + EtagReader { inner, bytes_tx, md5_rx } } - pub fn etag(self) -> String { - hex_simd::encode_to_string(self.md5.finalize(), hex_simd::AsciiCase::Lower) + pub async fn etag(self) -> String { + drop(self.inner); + drop(self.bytes_tx); + self.md5_rx.await.unwrap() } } impl AsyncRead for EtagReader { + #[tracing::instrument(level = "debug", skip_all)] fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { - let befor_size = buf.filled().len(); - - match Pin::new(&mut self.inner).poll_read(cx, buf) { - Poll::Ready(Ok(())) => { - if buf.filled().len() > befor_size { - let bytes = &buf.filled()[befor_size..]; - self.md5.update(bytes); - } - - Poll::Ready(Ok(())) + let poll = Pin::new(&mut self.inner).poll_read(cx, buf); + if let Poll::Ready(Ok(())) = &poll { + if buf.remaining() == 0 { + let bytes = buf.filled(); + let bytes = Bytes::copy_from_slice(bytes); + let tx = self.bytes_tx.clone(); + tokio::spawn(async move { + if let Err(e) = tx.send(bytes).await { + warn!("EtagReader send error: {:?}", e); + } + }); } - other => other, } + poll } } diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index c64bc189b..b35c6ad89 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -3802,7 +3802,7 @@ impl ObjectIO for SetDisks { error!("close_bitrot_writers err {:?}", err); } - let etag = etag_stream.etag(); + let etag = etag_stream.etag().await; //TODO: userDefined user_defined.insert("etag".to_owned(), etag.clone()); @@ -4393,7 +4393,7 @@ impl StorageAPI for SetDisks { error!("close_bitrot_writers err {:?}", err); } - let mut etag = etag_stream.etag(); + let mut etag = etag_stream.etag().await; if let Some(ref tag) = opts.preserve_etag { etag = tag.clone();