mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 16:28:15 +00:00
+13
-6
@@ -6,6 +6,7 @@ use crate::{
|
||||
};
|
||||
use blake2::Blake2b512;
|
||||
use blake2::Digest as _;
|
||||
use bytes::Bytes;
|
||||
use common::error::{Error, Result};
|
||||
use highway::{HighwayHash, HighwayHasher, Key};
|
||||
use lazy_static::lazy_static;
|
||||
@@ -533,20 +534,26 @@ impl Writer for BitrotFileWriter {
|
||||
self
|
||||
}
|
||||
|
||||
async fn write(&mut self, buf: &[u8]) -> Result<()> {
|
||||
async fn write(&mut self, buf: Bytes) -> Result<()> {
|
||||
if buf.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
self.hasher.reset();
|
||||
self.hasher.update(buf);
|
||||
let hash_bytes = self.hasher.clone().finalize();
|
||||
let mut hasher = self.hasher.clone();
|
||||
let h_buf = buf.clone();
|
||||
let hash_bytes = tokio::spawn(async move {
|
||||
hasher.reset();
|
||||
hasher.update(h_buf);
|
||||
hasher.finalize()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
if let Some(f) = self.inner.as_mut() {
|
||||
f.write_all(&hash_bytes).await?;
|
||||
f.write_all(buf).await?;
|
||||
f.write_all(&buf).await?;
|
||||
} else {
|
||||
self.inline_data.extend_from_slice(&hash_bytes);
|
||||
self.inline_data.extend_from_slice(buf);
|
||||
self.inline_data.extend_from_slice(&buf);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
|
||||
+16
-15
@@ -86,7 +86,6 @@ impl Erasure {
|
||||
}
|
||||
|
||||
self.buf.resize(new_len, 0u8);
|
||||
|
||||
match reader.read_exact(&mut self.buf).await {
|
||||
Ok(res) => res,
|
||||
Err(e) => {
|
||||
@@ -101,20 +100,22 @@ impl Erasure {
|
||||
}
|
||||
|
||||
self.encode_data(&self.buf, &mut blocks)?;
|
||||
let mut errs = Vec::new();
|
||||
|
||||
// TODO: 并发写入
|
||||
for (i, w_op) in writers.iter_mut().enumerate() {
|
||||
if let Some(w) = w_op {
|
||||
match w.write(blocks[i].as_ref()).await {
|
||||
Ok(_) => errs.push(None),
|
||||
Err(e) => errs.push(Some(e)),
|
||||
let write_futures = writers.iter_mut().enumerate().map(|(i, w_op)| {
|
||||
let i_inner = i.clone();
|
||||
let blocks_inner = blocks.clone();
|
||||
async move {
|
||||
if let Some(w) = w_op {
|
||||
match w.write(blocks_inner[i_inner].clone()).await {
|
||||
Ok(_) => None,
|
||||
Err(e) => Some(e),
|
||||
}
|
||||
} else {
|
||||
Some(Error::new(DiskError::DiskNotFound))
|
||||
}
|
||||
} else {
|
||||
errs.push(Some(Error::new(DiskError::DiskNotFound)));
|
||||
}
|
||||
}
|
||||
|
||||
});
|
||||
let errs = join_all(write_futures).await;
|
||||
let none_count = errs.iter().filter(|&x| x.is_none()).count();
|
||||
if none_count >= write_quorum {
|
||||
if total_size == 0 {
|
||||
@@ -472,7 +473,7 @@ impl Erasure {
|
||||
self.encoder.as_ref().unwrap().reconstruct(&mut bufs)?;
|
||||
}
|
||||
|
||||
let shards: Vec<Vec<u8>> = bufs.into_iter().flatten().collect::<Vec<_>>();
|
||||
let shards = bufs.into_iter().flatten().collect::<Vec<_>>();
|
||||
if shards.len() != self.parity_shards + self.data_shards {
|
||||
return Err(Error::from_string("can not reconstruct data"));
|
||||
}
|
||||
@@ -481,7 +482,7 @@ impl Erasure {
|
||||
if w.is_none() {
|
||||
continue;
|
||||
}
|
||||
match w.as_mut().unwrap().write(shards[i].as_ref()).await {
|
||||
match w.as_mut().unwrap().write(shards[i].clone().into()).await {
|
||||
Ok(_) => {}
|
||||
Err(e) => {
|
||||
info!("write failed, err: {:?}", e);
|
||||
@@ -501,7 +502,7 @@ impl Erasure {
|
||||
#[async_trait::async_trait]
|
||||
pub trait Writer {
|
||||
fn as_any(&self) -> &dyn Any;
|
||||
async fn write(&mut self, buf: &[u8]) -> Result<()>;
|
||||
async fn write(&mut self, buf: Bytes) -> Result<()>;
|
||||
async fn close(&mut self) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
+35
-15
@@ -1,3 +1,4 @@
|
||||
use bytes::Bytes;
|
||||
use futures::TryStreamExt;
|
||||
use md5::Digest;
|
||||
use md5::Md5;
|
||||
@@ -7,6 +8,7 @@ use std::task::Poll;
|
||||
use tokio::io::AsyncRead;
|
||||
use tokio::io::AsyncWrite;
|
||||
use tokio::io::ReadBuf;
|
||||
use tokio::sync::mpsc;
|
||||
use tokio::sync::oneshot;
|
||||
use tokio_util::io::ReaderStream;
|
||||
use tokio_util::io::StreamReader;
|
||||
@@ -125,33 +127,51 @@ impl AsyncRead for HttpFileReader {
|
||||
|
||||
pub struct EtagReader<R> {
|
||||
inner: R,
|
||||
md5: Md5,
|
||||
bytes_tx: mpsc::Sender<Bytes>,
|
||||
md5_rx: oneshot::Receiver<String>,
|
||||
}
|
||||
|
||||
impl<R> EtagReader<R> {
|
||||
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 {
|
||||
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<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<()>> {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
+24
-19
@@ -325,6 +325,7 @@ impl SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
let mut futures = Vec::with_capacity(disks.len());
|
||||
if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) {
|
||||
// TODO: 并发
|
||||
for (i, err) in errs.iter().enumerate() {
|
||||
@@ -334,26 +335,30 @@ impl SetDisks {
|
||||
|
||||
if let Some(disk) = disks[i].as_ref() {
|
||||
let fi = file_infos[i].clone();
|
||||
let _ = disk
|
||||
.delete_version(
|
||||
src_bucket,
|
||||
src_object,
|
||||
fi,
|
||||
false,
|
||||
DeleteOptions {
|
||||
undo_write: true,
|
||||
old_data_dir: data_dirs[i],
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
debug!("rename_data delete_version err {:?}", e);
|
||||
e
|
||||
});
|
||||
let old_data_dir = data_dirs[i];
|
||||
futures.push(async move {
|
||||
let _ = disk
|
||||
.delete_version(
|
||||
src_bucket,
|
||||
src_object,
|
||||
fi,
|
||||
false,
|
||||
DeleteOptions {
|
||||
undo_write: true,
|
||||
old_data_dir: old_data_dir,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
debug!("rename_data delete_version err {:?}", e);
|
||||
e
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
let _ = join_all(futures).await;
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
@@ -3793,7 +3798,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());
|
||||
@@ -4384,7 +4389,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();
|
||||
|
||||
Reference in New Issue
Block a user