support async calculate etag

Signed-off-by: junxiang Mu <1948535941@qq.com>
This commit is contained in:
junxiang Mu
2025-04-24 03:04:05 +00:00
parent e58eaacb20
commit 08d0021cd6
2 changed files with 37 additions and 17 deletions
+35 -15
View File
@@ -1,6 +1,8 @@
use bytes::Bytes;
use futures::TryStreamExt; use futures::TryStreamExt;
use md5::Digest; use md5::Digest;
use md5::Md5; use md5::Md5;
use tokio::sync::mpsc;
use std::pin::Pin; use std::pin::Pin;
use std::task::Context; use std::task::Context;
use std::task::Poll; use std::task::Poll;
@@ -125,33 +127,51 @@ impl AsyncRead for HttpFileReader {
pub struct EtagReader<R> { pub struct EtagReader<R> {
inner: R, inner: R,
md5: Md5, bytes_tx: mpsc::Sender<Bytes>,
md5_rx: oneshot::Receiver<String>,
} }
impl<R> EtagReader<R> { impl<R> EtagReader<R> {
pub fn new(inner: R) -> Self { pub fn new(inner: R) -> Self {
EtagReader { inner, md5: Md5::new() } let (bytes_tx, mut bytes_rx) = mpsc::channel::<Bytes>(8);
let (md5_tx, md5_rx) = oneshot::channel::<String>();
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 { pub async fn etag(self) -> String {
hex_simd::encode_to_string(self.md5.finalize(), hex_simd::AsciiCase::Lower) drop(self.inner);
drop(self.bytes_tx);
self.md5_rx.await.unwrap()
} }
} }
impl<R: AsyncRead + Unpin> AsyncRead for EtagReader<R> { impl<R: AsyncRead + Unpin> AsyncRead for EtagReader<R> {
#[tracing::instrument(level = "debug", skip_all)]
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<tokio::io::Result<()>> { fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<tokio::io::Result<()>> {
let befor_size = buf.filled().len(); let poll = Pin::new(&mut self.inner).poll_read(cx, buf);
if let Poll::Ready(Ok(())) = &poll {
match Pin::new(&mut self.inner).poll_read(cx, buf) { if buf.remaining() == 0 {
Poll::Ready(Ok(())) => { let bytes = buf.filled();
if buf.filled().len() > befor_size { let bytes = Bytes::copy_from_slice(bytes);
let bytes = &buf.filled()[befor_size..]; let tx = self.bytes_tx.clone();
self.md5.update(bytes); tokio::spawn(async move {
} if let Err(e) = tx.send(bytes).await {
warn!("EtagReader send error: {:?}", e);
Poll::Ready(Ok(())) }
});
} }
other => other,
} }
poll
} }
} }
+2 -2
View File
@@ -3802,7 +3802,7 @@ impl ObjectIO for SetDisks {
error!("close_bitrot_writers err {:?}", err); error!("close_bitrot_writers err {:?}", err);
} }
let etag = etag_stream.etag(); let etag = etag_stream.etag().await;
//TODO: userDefined //TODO: userDefined
user_defined.insert("etag".to_owned(), etag.clone()); user_defined.insert("etag".to_owned(), etag.clone());
@@ -4393,7 +4393,7 @@ impl StorageAPI for SetDisks {
error!("close_bitrot_writers err {:?}", err); 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 { if let Some(ref tag) = opts.preserve_etag {
etag = tag.clone(); etag = tag.clone();