diff --git a/ecstore/src/cache_value/metacache_set.rs b/ecstore/src/cache_value/metacache_set.rs index 79601ac87..4e2bbe579 100644 --- a/ecstore/src/cache_value/metacache_set.rs +++ b/ecstore/src/cache_value/metacache_set.rs @@ -12,6 +12,7 @@ use tokio::{ use crate::{ disk::{DiskAPI, DiskStore, MetaCacheEntries, MetaCacheEntry, WalkDirOptions}, error::{Error, Result}, + io::Writer, }; type AgreedFn = Box Pin + Send>> + Send + 'static>; @@ -77,16 +78,19 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - match disk .as_ref() .unwrap() - .walk_dir(WalkDirOptions { - bucket: opts_clone.bucket.clone(), - base_dir: opts_clone.path.clone(), - recursive: opts_clone.recursice, - report_notfound: opts_clone.report_not_found, - filter_prefix: opts_clone.filter_prefix.clone(), - forward_to: opts_clone.forward_to.clone(), - limit: opts_clone.per_disk_limit, - ..Default::default() - }) + .walk_dir( + WalkDirOptions { + bucket: opts_clone.bucket.clone(), + base_dir: opts_clone.path.clone(), + recursive: opts_clone.recursice, + report_notfound: opts_clone.report_not_found, + filter_prefix: opts_clone.filter_prefix.clone(), + forward_to: opts_clone.forward_to.clone(), + limit: opts_clone.per_disk_limit, + ..Default::default() + }, + Writer::NotUse, + ) .await { Ok(r) => { @@ -115,16 +119,19 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - match disk .as_ref() .unwrap() - .walk_dir(WalkDirOptions { - bucket: opts_clone.bucket.clone(), - base_dir: opts_clone.path.clone(), - recursive: opts_clone.recursice, - report_notfound: opts_clone.report_not_found, - filter_prefix: opts_clone.filter_prefix.clone(), - forward_to: opts_clone.forward_to.clone(), - limit: opts_clone.per_disk_limit, - ..Default::default() - }) + .walk_dir( + WalkDirOptions { + bucket: opts_clone.bucket.clone(), + base_dir: opts_clone.path.clone(), + recursive: opts_clone.recursice, + report_notfound: opts_clone.report_not_found, + filter_prefix: opts_clone.filter_prefix.clone(), + forward_to: opts_clone.forward_to.clone(), + limit: opts_clone.per_disk_limit, + ..Default::default() + }, + Writer::NotUse, + ) .await { Ok(r) => { diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 8cca640b6..268c672b5 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -1291,10 +1291,18 @@ impl DiskAPI for LocalDisk { Ok(entries) } - // TODO: io.writer - async fn walk_dir(&self, opts: WalkDirOptions) -> Result> { + // FIXME: TODO: io.writer + async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result> { // warn!("walk_dir opts {:?}", &opts); + let volume_dir = self.get_bucket_path(&opts.bucket)?; + + if !skip_access_checks(&opts.bucket) { + if let Err(e) = access(&volume_dir).await { + return Err(convert_access_error(e, DiskError::VolumeAccessDenied)); + } + } + let mut metas = Vec::new(); if opts.base_dir.ends_with(SLASH_SEPARATOR) { diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index dd1fc6c09..c7c5a1003 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -22,6 +22,7 @@ use crate::{ data_usage_cache::{DataUsageCache, DataUsageEntry}, heal_commands::{HealScanMode, HealingTracker}, }, + io, store_api::{FileInfo, RawFileInfo}, }; use endpoint::Endpoint; @@ -208,10 +209,10 @@ impl DiskAPI for Disk { } } - async fn walk_dir(&self, opts: WalkDirOptions) -> Result> { + async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result> { match self { - Disk::Local(local_disk) => local_disk.walk_dir(opts).await, - Disk::Remote(remote_disk) => remote_disk.walk_dir(opts).await, + Disk::Local(local_disk) => local_disk.walk_dir(opts, wr).await, + Disk::Remote(remote_disk) => remote_disk.walk_dir(opts, wr).await, } } @@ -403,7 +404,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) -> Result>; + async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result>; // Metadata operations async fn delete_version( @@ -599,10 +600,10 @@ pub struct MetaCacheEntry { pub metadata: Vec, // cached contains the metadata if decoded. - cached: Option, + pub cached: Option, // Indicates the entry can be reused and only one reference to metadata is expected. - _reusable: bool, + pub reusable: bool, } impl MetaCacheEntry { @@ -616,6 +617,7 @@ impl MetaCacheEntry { Ok(wr) } + pub fn is_dir(&self) -> bool { self.metadata.is_empty() && self.name.ends_with('/') } @@ -838,7 +840,7 @@ impl MetaCacheEntries { meta_ver: selected.as_ref().unwrap().cached.as_ref().unwrap().meta_ver, ..Default::default() }), - _reusable: true, + reusable: true, ..Default::default() }); diff --git a/ecstore/src/disk/remote.rs b/ecstore/src/disk/remote.rs index 618365f6d..3259f9414 100644 --- a/ecstore/src/disk/remote.rs +++ b/ecstore/src/disk/remote.rs @@ -346,7 +346,7 @@ impl DiskAPI for RemoteDisk { Ok(response.volumes) } - async fn walk_dir(&self, opts: WalkDirOptions) -> Result> { + async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> 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 new file mode 100644 index 000000000..ca2c3904e --- /dev/null +++ b/ecstore/src/io.rs @@ -0,0 +1,68 @@ +use std::pin::Pin; +use std::task::{Context, Poll}; +use tokio::fs::File; +use tokio::io::{self, AsyncRead, AsyncWrite, ReadBuf}; + +#[derive(Default)] +pub enum Reader { + #[default] + NotUse, + File(File), +} + +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::NotUse => Poll::Ready(Ok(())), + } + } +} + +#[derive(Default)] +pub enum Writer { + #[default] + NotUse, + File(File), +} + +impl AsyncWrite for Writer { + fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll> { + match self.get_mut() { + Writer::File(file) => { + // Create a pinned reference from the file + let file = Pin::new(file); + file.poll_write(cx, buf) + } + Writer::NotUse => Poll::Ready(Ok(0)), + } + } + + fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + match self.get_mut() { + Writer::File(file) => { + let file = Pin::new(file); + file.poll_flush(cx) + } + Writer::NotUse => Poll::Ready(Ok(())), + } + } + + fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + match self.get_mut() { + Writer::File(file) => { + let file = Pin::new(file); + file.poll_shutdown(cx) + } + Writer::NotUse => Poll::Ready(Ok(())), + } + } +} + +// #[tokio::test] +// async fn test_reader{ + +// } diff --git a/ecstore/src/lib.rs b/ecstore/src/lib.rs index 42623a660..8aa575a33 100644 --- a/ecstore/src/lib.rs +++ b/ecstore/src/lib.rs @@ -13,6 +13,8 @@ mod file_meta; pub mod file_meta_inline; pub mod global; pub mod heal; +pub mod io; +pub mod metacache; pub mod metrics_realtime; pub mod notification_sys; pub mod peer; diff --git a/ecstore/src/metacache/mod.rs b/ecstore/src/metacache/mod.rs new file mode 100644 index 000000000..d3baa8178 --- /dev/null +++ b/ecstore/src/metacache/mod.rs @@ -0,0 +1 @@ +pub mod writer; diff --git a/ecstore/src/metacache/writer.rs b/ecstore/src/metacache/writer.rs new file mode 100644 index 000000000..0be530eb8 --- /dev/null +++ b/ecstore/src/metacache/writer.rs @@ -0,0 +1,206 @@ +use std::io::ErrorKind; +use std::io::Read; +use std::io::Write; +use std::str::from_utf8; + +use crate::disk::MetaCacheEntry; +use crate::error::Error; +use crate::error::Result; + +const METACACHE_STREAM_VERSION: u8 = 2; + +pub struct MetacacheWriter { + wr: W, + buf: Vec, + created: bool, +} + +impl MetacacheWriter { + pub fn new(wr: W, block_size: usize) -> Self { + Self { + wr, + buf: Vec::with_capacity(block_size), + created: false, + } + } + + async fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> { + if objs.is_empty() { + return Ok(()); + } + + if !self.created { + rmp::encode::write_u8(&mut self.wr, METACACHE_STREAM_VERSION)?; + self.created = false; + } + + for obj in objs.iter() { + if obj.name.is_empty() { + return Err(Error::msg("metacacheWriter: no name")); + } + + rmp::encode::write_bool(&mut self.wr, true)?; + + rmp::encode::write_str(&mut self.wr, &obj.name)?; + + rmp::encode::write_bin(&mut self.wr, &obj.metadata)?; + } + + Ok(()) + } + + async fn close(&mut self) -> Result<()> { + rmp::encode::write_bool(&mut self.wr, false)?; + + self.wr.flush()?; + Ok(()) + } +} + +pub struct MetacacheReader { + rd: R, + init: bool, + err: Option, + buf: Vec, +} + +impl MetacacheReader { + pub fn new(rd: R) -> Self { + Self { + rd, + init: false, + err: None, + buf: Vec::new(), + } + } + + pub fn check_init(&mut self) { + if !self.init { + let ver = match rmp::decode::read_u8(&mut self.rd) { + Ok(res) => res, + Err(err) => { + self.err = Some(Error::msg(err.to_string())); + 0 + } + }; + match ver { + 1 | 2 => (), + _ => { + self.err = Some(Error::msg("invalid version")); + } + } + + self.init = true; + } + } + + pub fn peek(&mut self) -> Result { + self.check_init(); + + if let Some(err) = &self.err { + return Err(err.clone()); + } + + match rmp::decode::read_bool(&mut self.rd) { + Ok(res) => { + if !res { + self.err = Some(Error::new(std::io::Error::from(ErrorKind::UnexpectedEof))); + return Err(Error::new(std::io::Error::from(ErrorKind::UnexpectedEof))); + } + } + Err(err) => { + self.err = Some(Error::msg(err.to_string())); + return Err(Error::new(err)); + } + }; + + let l = match rmp::decode::read_str_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); + 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 = 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 metadata = self.buf.clone(); + + Ok(MetaCacheEntry { + name, + metadata, + cached: None, + reusable: false, + }) + } +} + +#[tokio::test] +async fn test_writer() { + use std::fs::File; + use std::fs::OpenOptions; + + let file_path = "./test_writer.txt"; + let f = OpenOptions::new() + .create(true) + .read(true) + .write(true) + .truncate(true) + .open(file_path) + .unwrap(); + + // let wr = Writer::File(f); + + let mut w = MetacacheWriter::new(f, 1024); + + let mut objs = Vec::new(); + for i in 0..10 { + objs.push(MetaCacheEntry { + name: format!("item{}", i), + metadata: vec![0u8, 10], + cached: None, + reusable: false, + }); + } + + w.write(&objs).await.unwrap(); + w.close().await.unwrap(); + + let nf = File::open(file_path).unwrap(); + + let meta = nf.metadata().unwrap(); + + println!("{}", meta.len()); +} diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 4061024bc..4a6b430e4 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -6,7 +6,6 @@ use std::{ time::Duration, }; -use crate::heal::heal_ops::{HealEntryFn, HealSequence}; use crate::{ bitrot::{bitrot_verify, close_bitrot_writers, new_bitrot_filereader, new_bitrot_filewriter, BitrotFileWriter}, cache_value::metacache_set::{list_path_raw, ListPathRawOptions}, @@ -56,6 +55,10 @@ use crate::{ heal::data_scanner::{globalHealConfig, HEAL_DELETE_DANGLING}, store_api::ListObjectVersionsInfo, }; +use crate::{ + heal::heal_ops::{HealEntryFn, HealSequence}, + io::Writer, +}; use futures::future::join_all; use glob::Pattern; use http::HeaderMap; @@ -1333,7 +1336,7 @@ impl SetDisks { let disk = disk.as_ref().unwrap(); let opts = opts.clone(); // let mut wr = &mut wr; - futures.push(disk.walk_dir(opts)); + futures.push(disk.walk_dir(opts, Writer::NotUse)); // tokio::spawn(async move { disk.walk_dir(opts, wr).await }); } diff --git a/rustfs/src/grpc.rs b/rustfs/src/grpc.rs index 25c3391d9..16c7c588c 100644 --- a/rustfs/src/grpc.rs +++ b/rustfs/src/grpc.rs @@ -743,7 +743,7 @@ impl Node for NodeService { })); } }; - match disk.walk_dir(opts).await { + match disk.walk_dir(opts, ecstore::io::Writer::NotUse).await { Ok(entries) => { let entries = entries .into_iter()