From 99c229524d0d3634513e50a85bf6488a4b5461e7 Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 16 Dec 2024 17:45:28 +0800 Subject: [PATCH] test metawrite --- ecstore/src/cache_value/metacache_set.rs | 4 +- ecstore/src/disk/local.rs | 66 ++++--- ecstore/src/disk/mod.rs | 6 +- ecstore/src/disk/remote.rs | 7 +- ecstore/src/io.rs | 15 +- ecstore/src/metacache/writer.rs | 229 +++++++++++++++-------- ecstore/src/set_disk.rs | 2 +- rustfs/src/grpc.rs | 2 +- 8 files changed, 210 insertions(+), 121 deletions(-) diff --git a/ecstore/src/cache_value/metacache_set.rs b/ecstore/src/cache_value/metacache_set.rs index 4e2bbe579..a7ed75af0 100644 --- a/ecstore/src/cache_value/metacache_set.rs +++ b/ecstore/src/cache_value/metacache_set.rs @@ -89,7 +89,7 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - limit: opts_clone.per_disk_limit, ..Default::default() }, - Writer::NotUse, + &mut Writer::NotUse, ) .await { @@ -130,7 +130,7 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - limit: opts_clone.per_disk_limit, ..Default::default() }, - Writer::NotUse, + &mut Writer::NotUse, ) .await { diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index aac07a218..5a8979bb4 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -59,10 +59,10 @@ use std::{ }; use time::OffsetDateTime; use tokio::fs::{self, File}; -use tokio::io::{AsyncReadExt, AsyncWriteExt, ErrorKind}; +use tokio::io::{AsyncReadExt, AsyncWrite, AsyncWriteExt, ErrorKind}; use tokio::sync::mpsc::Sender; use tokio::sync::RwLock; -use tracing::{info, warn}; +use tracing::{error, info, warn}; use uuid::Uuid; #[derive(Debug)] @@ -716,7 +716,7 @@ impl LocalDisk { bitrot_verify(&mut Cursor::new(data), n, part_size, algo, sum.to_vec(), shard_size) } - async fn scan_dir( + async fn scan_dir( &self, current: &mut String, opts: &WalkDirOptions, @@ -810,7 +810,8 @@ impl LocalDisk { name, metadata, ..Default::default() - })?; + }) + .await?; *objs_returned += 1; return Ok(()); } @@ -850,7 +851,8 @@ impl LocalDisk { out.write_obj(&MetaCacheEntry { name: pop, ..Default::default() - })?; + }) + .await?; if opts.recursive { if let Err(er) = Box::pin(self.scan_dir(current, opts, out, objs_returned)).await { @@ -886,7 +888,7 @@ impl LocalDisk { meta.metadata = res; - out.write_obj(&meta)?; + out.write_obj(&meta).await?; *objs_returned += 1; } Err(err) => { @@ -914,7 +916,8 @@ impl LocalDisk { out.write_obj(&MetaCacheEntry { name: dir.clone(), ..Default::default() - })?; + }) + .await?; *objs_returned += 1; if opts.recursive { @@ -1511,7 +1514,7 @@ impl DiskAPI for LocalDisk { // FIXME: TODO: io.writer TODO cancel #[tracing::instrument(level = "debug", skip(self, wr))] - async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result> { + async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result> { // warn!("walk_dir opts {:?}", &opts); let volume_dir = self.get_bucket_path(&opts.bucket)?; @@ -1524,7 +1527,7 @@ impl DiskAPI for LocalDisk { let mut wr = wr; - let mut out = MetacacheWriter::new(AsyncToSync::new_writer(&mut wr)); + let mut out = MetacacheWriter::new(&mut wr); let mut objs_returned = 0; @@ -1539,7 +1542,7 @@ impl DiskAPI for LocalDisk { metadata: data, ..Default::default() }; - out.write_obj(&meta)?; + out.write_obj(&meta).await?; objs_returned += 1; } } @@ -2330,6 +2333,10 @@ mod test { use utils::fs::O_RDWR; use utils::fs::O_TRUNC; + use crate::io::VecAsyncReader; + use crate::io::VecAsyncWriter; + use crate::metacache::writer::MetacacheReader; + use super::*; #[tokio::test] @@ -2425,25 +2432,32 @@ mod test { } }; - let f = match open_file("./testfile.txt", O_CREATE | O_RDWR | O_TRUNC).await { - Ok(res) => res, - Err(err) => { - println!("openfile err {:?}", err); - return; + let (rd, mut wr) = tokio::io::duplex(64); + + // let mut wr = VecAsyncWriter::new(Vec::new()); + + let job = tokio::spawn(async move { + let opts = WalkDirOptions { + bucket: "dada".to_owned(), + recursive: true, + ..Default::default() + }; + if let Err(err) = disk.walk_dir(opts, &mut wr).await { + println!("walk_dir err {:?}", err); } - }; + }); - let buf = BufWriter::new(Vec::new()); + // let rd = VecAsyncReader::new(wr.get_buffer().to_vec()); - let opts = WalkDirOptions { - bucket: "dada".to_owned(), - recursive: true, - ..Default::default() - }; - if let Err(err) = disk.walk_dir(opts, crate::io::Writer::File(f)).await { - println!("walk_dir err {:?}", err); - } + let rd_job = tokio::spawn(async move { + let mut mrd = MetacacheReader::new(rd); - MetacacheReader::new() + while let Some(info) = mrd.peek().await.unwrap_or_default() { + println!("{:?}", info) + } + }); + + job.await.unwrap(); + rd_job.await.unwrap(); } } diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 5fec3bc29..c47dd9afe 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -46,7 +46,7 @@ use std::{ use time::OffsetDateTime; use tokio::{ fs::File, - io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt}, + io::{AsyncReadExt, AsyncSeekExt, AsyncWrite, AsyncWriteExt}, sync::mpsc::{self, Sender}, }; use tokio_stream::wrappers::ReceiverStream; @@ -210,7 +210,7 @@ impl DiskAPI for Disk { } } - async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result> { + async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result> { match self { Disk::Local(local_disk) => local_disk.walk_dir(opts, wr).await, Disk::Remote(remote_disk) => remote_disk.walk_dir(opts, wr).await, @@ -406,7 +406,7 @@ pub trait DiskAPI: Debug + Send + Sync + 'static { async fn delete_volume(&self, volume: &str) -> Result<()>; // 并发边读边写 TODO: wr io.Writer - async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result>; + async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result>; // Metadata operations async fn delete_version( diff --git a/ecstore/src/disk/remote.rs b/ecstore/src/disk/remote.rs index 6f5534c79..a68d6d827 100644 --- a/ecstore/src/disk/remote.rs +++ b/ecstore/src/disk/remote.rs @@ -10,7 +10,10 @@ use protos::{ StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WalkDirRequest, WriteAllRequest, WriteMetadataRequest, }, }; -use tokio::sync::mpsc::{self, Sender}; +use tokio::{ + io::AsyncWrite, + sync::mpsc::{self, Sender}, +}; use tokio_stream::{wrappers::ReceiverStream, StreamExt}; use tonic::Request; use tracing::info; @@ -347,7 +350,7 @@ impl DiskAPI for RemoteDisk { Ok(response.volumes) } - async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result> { + async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result> { info!("walk_dir"); let walk_dir_options = serde_json::to_string(&opts)?; let mut client = node_service_time_out_client(&self.addr) diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index a4d9c3cc1..a6b182ce8 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -3,8 +3,6 @@ use std::io::Write; use std::pin::Pin; use std::task::{Context, Poll}; use tokio::fs::File; -use tokio::io::BufReader; -use tokio::io::BufWriter; use tokio::io::{self, AsyncRead, AsyncWrite, ReadBuf}; #[derive(Default)] @@ -12,19 +10,14 @@ pub enum Reader { #[default] NotUse, File(File), - Buffer(BufReader>), + Buffer(VecAsyncReader), } impl AsyncRead for Reader { fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { match self.get_mut() { - Reader::File(file) => { - let file = Pin::new(file); - file.poll_read(cx, buf) - } - Reader::Buffer(buffer) => { - todo!() - } + Reader::File(file) => Pin::new(file).poll_read(cx, buf), + Reader::Buffer(buffer) => Pin::new(buffer).poll_read(cx, buf), Reader::NotUse => Poll::Ready(Ok(())), } } @@ -35,7 +28,7 @@ pub enum Writer { #[default] NotUse, File(File), - Buffer(BufWriter>), + Buffer(VecAsyncWriter), } impl AsyncWrite for Writer { diff --git a/ecstore/src/metacache/writer.rs b/ecstore/src/metacache/writer.rs index cc23cf6ec..ddc55f913 100644 --- a/ecstore/src/metacache/writer.rs +++ b/ecstore/src/metacache/writer.rs @@ -1,9 +1,16 @@ use crate::disk::MetaCacheEntry; use crate::error::Error; use crate::error::Result; +use rmp::decode::RmpRead; +use rmp::encode::RmpWrite; +use rmp::Marker; use std::io::Read; use std::io::Write; use std::str::from_utf8; +use tokio::io::AsyncRead; +use tokio::io::AsyncReadExt; +use tokio::io::AsyncWrite; +use tokio::io::AsyncWriteExt; // use std::sync::Arc; // use tokio::sync::mpsc; // use tokio::sync::mpsc::Sender; @@ -16,50 +23,64 @@ pub struct MetacacheWriter { wr: W, created: bool, // err: Option, + buf: Vec, } -impl MetacacheWriter { +impl MetacacheWriter { pub fn new(wr: W) -> Self { Self { wr, created: false, // err: None, + buf: Vec::new(), } } - pub fn init(&mut self) -> Result<()> { + pub async fn flush(&mut self) -> Result<()> { + self.wr.write_all(&self.buf).await?; + self.buf.clear(); + + Ok(()) + } + + pub async fn init(&mut self) -> Result<()> { if !self.created { - rmp::encode::write_u8(&mut self.wr, METACACHE_STREAM_VERSION)?; + rmp::encode::write_u8(&mut self.buf, METACACHE_STREAM_VERSION).map_err(|e| Error::msg(format!("{:?}", e)))?; + self.flush().await?; self.created = true; } Ok(()) } - pub fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> { + pub async fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> { if objs.is_empty() { return Ok(()); } - self.init()?; + self.init().await?; for obj in objs.iter() { if obj.name.is_empty() { return Err(Error::msg("metacacheWriter: no name")); } - self.write_obj(obj)?; + self.write_obj(obj).await?; } Ok(()) } - pub fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> { - self.init()?; - rmp::encode::write_bool(&mut self.wr, true)?; + pub async fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> { + self.init().await?; - rmp::encode::write_str(&mut self.wr, &obj.name)?; + rmp::encode::write_bool(&mut self.buf, true).map_err(|e| Error::msg(format!("{:?}", e)))?; + + rmp::encode::write_str(&mut self.buf, &obj.name).map_err(|e| Error::msg(format!("{:?}", e)))?; + + rmp::encode::write_bin(&mut self.buf, &obj.metadata).map_err(|e| Error::msg(format!("{:?}", e)))?; + + self.flush().await?; - rmp::encode::write_bin(&mut self.wr, &obj.metadata)?; Ok(()) } @@ -96,9 +117,9 @@ impl MetacacheWriter { // Ok(sender) // } - pub fn close(&mut self) -> Result<()> { - rmp::encode::write_bool(&mut self.wr, false)?; - self.wr.flush()?; + pub async fn close(&mut self) -> Result<()> { + rmp::encode::write_bool(&mut self.buf, false).map_err(|e| Error::msg(format!("{:?}", e)))?; + self.flush().await?; Ok(()) } } @@ -108,34 +129,60 @@ pub struct MetacacheReader { init: bool, err: Option, buf: Vec, + offset: usize, } -impl MetacacheReader { +impl MetacacheReader { pub fn new(rd: R) -> Self { Self { rd, init: false, err: None, buf: Vec::new(), + offset: 0, } } - fn check_init(&mut self) { + pub async fn read_more(&mut self, read_size: usize) -> Result<&[u8]> { + let ext_size = read_size + self.offset; + + let extra = ext_size - self.offset; + if self.buf.capacity() >= ext_size { + // Extend the buffer if we have enough space. + self.buf.resize(ext_size, 0); + } else { + self.buf.extend(vec![0u8; extra]); + } + + let pref = self.offset; + + self.rd.read_exact(&mut self.buf[pref..ext_size]).await?; + + self.offset += read_size; + + let data = &self.buf[pref..ext_size]; + + println!("pref {} offset {},ext_size {}, data {:?}", pref, self.offset, ext_size, &data); + + Ok(data) + } + + fn reset(&mut self) { + self.buf.clear(); + self.offset = 0; + } + + async fn check_init(&mut self) -> Result<()> { if !self.init { - // let mut buf = match self.read_buf(1).await { - // Ok(res) => res, - // Err(err) => { - // self.err = Some(Error::msg(err.to_string())); - // return; - // } - // }; - let ver = match rmp::decode::read_u8(&mut self.rd) { + let ver = match rmp::decode::read_u8(&mut self.read_more(2).await?) { Ok(res) => res, Err(err) => { - self.err = Some(Error::msg(err.to_string())); + self.err = Some(Error::msg(format!("{:?}", err))); 0 } }; + + println!("ver {}", ver); match ver { 1 | 2 => (), _ => { @@ -145,70 +192,99 @@ impl MetacacheReader { self.init = true; } + Ok(()) } - pub fn peek(&mut self) -> Result> { - self.check_init(); + async fn read_str_len(&mut self) -> Result { + let mark = match rmp::decode::read_marker(&mut self.read_more(1).await?) { + Ok(res) => res, + Err(err) => { + let serr = format!("{:?}", err); + self.err = Some(Error::msg(&serr)); + return Err(Error::msg(&serr)); + } + }; + + match mark { + Marker::FixStr(size) => Ok(u32::from(size)), + Marker::Str8 => Ok(u32::from(self.read_u8().await?)), + Marker::Str16 => Ok(u32::from(self.read_u16().await?)), + Marker::Str32 => Ok(self.read_u32().await?), + _ => Err(Error::msg("str marker err")), + } + } + + async fn read_bin_len(&mut self) -> Result { + let mark = match rmp::decode::read_marker(&mut self.read_more(1).await?) { + Ok(res) => res, + Err(err) => { + let serr = format!("{:?}", err); + self.err = Some(Error::msg(&serr)); + return Err(Error::msg(&serr)); + } + }; + + match mark { + Marker::Bin8 => Ok(u32::from(self.read_u8().await?)), + Marker::Bin16 => Ok(u32::from(self.read_u16().await?)), + Marker::Bin32 => Ok(self.read_u32().await?), + _ => Err(Error::msg("bin marker err")), + } + } + + async fn read_u8(&mut self) -> Result { + let a = self.read_more(1).await?; + + Ok(a[0]) + } + + async fn read_u16(&mut self) -> Result { + rmp::decode::read_u16(&mut self.read_more(2).await?).map_err(|e| Error::msg(format!("{:?}", e))) + } + + async fn read_u32(&mut self) -> Result { + rmp::decode::read_u32(&mut self.read_more(4).await?).map_err(|e| Error::msg(format!("{:?}", e))) + } + + pub async fn peek(&mut self) -> Result> { + self.check_init().await; if let Some(err) = &self.err { return Err(err.clone()); } - match rmp::decode::read_bool(&mut self.rd) { + match rmp::decode::read_bool(&mut self.read_more(1).await?) { Ok(res) => { if !res { return Ok(None); } } Err(err) => { - self.err = Some(Error::msg(err.to_string())); - return Err(Error::new(err)); + let serr = format!("{:?}", err); + self.err = Some(Error::msg(&serr)); + return Err(Error::msg(&serr)); } }; - let l = match rmp::decode::read_str_len(&mut self.rd) { - Ok(res) => res, + let l = self.read_str_len().await?; + + let buf = self.read_more(l as usize).await?; + let name_buf = buf.to_vec(); + let name = match from_utf8(&name_buf) { + Ok(decoded) => decoded.to_owned(), Err(err) => { self.err = Some(Error::msg(err.to_string())); - return Err(Error::new(err)); + return Err(Error::msg(err.to_string())); } }; - self.buf.resize(l as usize, 0); - let name = match self.rd.read_exact(&mut self.buf) { - Ok(()) => { - let name_buf = self.buf.to_vec(); - match from_utf8(&name_buf) { - Ok(decoded) => Ok(decoded.to_owned()), - Err(err) => { - self.err = Some(Error::msg(err.to_string())); - Err(Error::msg(err.to_string())) - } - } - } - Err(err) => { - self.err = Some(Error::msg(err.to_string())); - Err(Error::msg(err.to_string())) - } - }?; + let l = self.read_bin_len().await?; - let l = match rmp::decode::read_bin_len(&mut self.rd) { - Ok(res) => res, - Err(err) => { - self.err = Some(Error::msg(err.to_string())); - return Err(Error::new(err)); - } - }; - self.buf.resize(l as usize, 0); - match self.rd.read_exact(&mut self.buf) { - Ok(res) => res, - Err(err) => { - self.err = Some(Error::msg(err.to_string())); - return Err(Error::new(err)); - } - }; + let buf = self.read_more(l as usize).await?; - let metadata = self.buf.clone(); + let metadata = buf.to_vec(); + + self.reset(); Ok(Some(MetaCacheEntry { name, @@ -218,11 +294,11 @@ impl MetacacheReader { })) } - pub fn read_all(&mut self) -> Result> { + pub async fn read_all(&mut self) -> Result> { let mut ret = Vec::new(); loop { - if let Some(entry) = self.peek()? { + if let Some(entry) = self.peek().await? { ret.push(entry); continue; } @@ -236,13 +312,12 @@ impl MetacacheReader { #[tokio::test] async fn test_writer() { - use crate::io::AsyncToSync; use crate::io::VecAsyncReader; use crate::io::VecAsyncWriter; let mut f = VecAsyncWriter::new(Vec::new()); - let mut w = MetacacheWriter::new(AsyncToSync::new_writer(&mut f)); + let mut w = MetacacheWriter::new(&mut f); let mut objs = Vec::new(); for i in 0..10 { @@ -256,14 +331,18 @@ async fn test_writer() { objs.push(info); } - w.write(&objs).unwrap(); + w.write(&objs).await.unwrap(); - w.close().unwrap(); + w.close().await.unwrap(); - let nf = VecAsyncReader::new(f.get_buffer().to_vec()); + let data = f.get_buffer().to_vec(); - let mut r = MetacacheReader::new(AsyncToSync::new_reader(nf)); - let nobjs = r.read_all().unwrap(); + println!("data len {}", data.len()); + + let nf = VecAsyncReader::new(data); + + let mut r = MetacacheReader::new(nf); + let nobjs = r.read_all().await.unwrap(); for info in nobjs.iter() { println!("new {:?}", &info); diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 621ab8e6e..a7d1ab0ac 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -1330,7 +1330,7 @@ impl SetDisks { let opts = opts.clone(); futures.push(async move { if let Some(disk) = disk { - disk.walk_dir(opts, Writer::NotUse).await + disk.walk_dir(opts, &mut Writer::NotUse).await } else { Err(Error::new(DiskError::DiskNotFound)) } diff --git a/rustfs/src/grpc.rs b/rustfs/src/grpc.rs index c529ac8ca..a2be08bac 100644 --- a/rustfs/src/grpc.rs +++ b/rustfs/src/grpc.rs @@ -743,7 +743,7 @@ impl Node for NodeService { })); } }; - match disk.walk_dir(opts, ecstore::io::Writer::NotUse).await { + match disk.walk_dir(opts, &mut ecstore::io::Writer::NotUse).await { Ok(entries) => { let entries = entries .into_iter()