From 4195371a1ebae955c0d5a9abf9d61e7018276c3e Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 17 Dec 2024 14:28:37 +0800 Subject: [PATCH] test_list_path_raw done --- ecstore/src/cache_value/metacache_set.rs | 340 ++++++++++++++--------- ecstore/src/disk/error.rs | 11 + ecstore/src/disk/local.rs | 47 ---- ecstore/src/disk/mod.rs | 7 +- ecstore/src/endpoints.rs | 3 +- ecstore/src/error.rs | 7 +- ecstore/src/file_meta.rs | 4 - ecstore/src/metacache/writer.rs | 51 +++- ecstore/src/store_list_objects.rs | 129 +++++++++ 9 files changed, 401 insertions(+), 198 deletions(-) diff --git a/ecstore/src/cache_value/metacache_set.rs b/ecstore/src/cache_value/metacache_set.rs index a7ed75af0..2f719485d 100644 --- a/ecstore/src/cache_value/metacache_set.rs +++ b/ecstore/src/cache_value/metacache_set.rs @@ -1,5 +1,6 @@ use std::{future::Future, pin::Pin, sync::Arc}; +use futures::{future::join_all, join}; use tokio::{ spawn, sync::{ @@ -8,11 +9,16 @@ use tokio::{ RwLock, }, }; +use tracing::error; use crate::{ - disk::{DiskAPI, DiskStore, MetaCacheEntries, MetaCacheEntry, WalkDirOptions}, + disk::{ + error::{is_err_eof, is_err_file_not_found, is_err_volume_not_found, DiskError}, + DiskAPI, DiskStore, MetaCacheEntries, MetaCacheEntry, WalkDirOptions, + }, error::{Error, Result}, io::Writer, + metacache::writer::MetacacheReader, }; type AgreedFn = Box Pin + Send>> + Send + 'static>; @@ -62,44 +68,40 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - return Err(Error::from_string("list_path_raw: 0 drives provided")); } + let mut jobs: Vec>> = Vec::new(); let mut readers = Vec::with_capacity(opts.disks.len()); let fds = Arc::new(RwLock::new(opts.fallback_disks.clone())); for disk in opts.disks.iter() { - let disk = disk.clone(); + let opdisk = disk.clone(); let opts_clone = opts.clone(); let fds_clone = fds.clone(); - let (m_tx, m_rx) = mpsc::channel::(100); - readers.push(m_rx); - spawn(async move { + // let (m_tx, m_rx) = mpsc::channel::(100); + // readers.push(m_rx); + let (rd, mut wr) = tokio::io::duplex(64); + readers.push(MetacacheReader::new(rd)); + jobs.push(spawn(async move { + let wakl_opts = 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() + }; + let mut need_fallback = false; - if disk.is_none() { - need_fallback = true; - } else { - 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() - }, - &mut Writer::NotUse, - ) - .await - { - Ok(r) => { - for v in r.iter() { - let _ = m_tx.send(v.to_owned()).await; - } + if let Some(disk) = opdisk { + match disk.walk_dir(wakl_opts, &mut wr).await { + Ok(_res) => {} + Err(err) => { + error!("walk dir err {:?}", &err); + need_fallback = true; } - Err(_) => need_fallback = true, } + } else { + need_fallback = true; } while need_fallback { @@ -108,133 +110,203 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - if fds_w.is_empty() { break None; } - let fd = fds_w.remove(0); - if fd.is_some() && fd.as_ref().unwrap().is_online().await { - break fd; + + if let Some(fd) = fds_w.remove(0) { + if fd.is_online().await { + break Some(fd); + } } }; - if f_disk.is_none() { + + if let Some(disk) = f_disk { + match disk + .as_ref() + .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() + }, + &mut wr, + ) + .await + { + Ok(_r) => { + need_fallback = false; + } + Err(err) => { + error!("walk dir2 err {:?}", &err); + break; + } + } + } else { break; } - 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() - }, - &mut Writer::NotUse, - ) - .await - { - Ok(r) => { - for v in r.iter() { - let _ = m_tx.send(v.to_owned()).await; - } - need_fallback = false; - } - Err(_) => break, - } } - }); + + Ok(()) + })); } - let errs: Vec> = vec![None; readers.len()]; - loop { - let mut current = MetaCacheEntry::default(); - let (mut at_eof, mut has_err, mut agree) = (0, 0, 0); - if rx.try_recv().is_ok() { - return Err(Error::from_string("canceled")); - } - let mut top_entries: Vec = Vec::with_capacity(readers.len()); - // top_entries.clear(); + let revjob = spawn(async move { + let mut errs: Vec> = vec![None; readers.len()]; + loop { + let mut current = MetaCacheEntry::default(); - for (i, r) in readers.iter_mut().enumerate() { - if errs[i].is_none() { - has_err += 1; - continue; + if rx.try_recv().is_ok() { + return Err(Error::from_string("canceled")); } - let entry = match r.recv().await { - Some(entry) => entry, - None => { - at_eof += 1; + let mut top_entries: Vec> = vec![None; readers.len()]; + + let mut at_eof = 0; + let mut fnf = 0; + let mut vnf = 0; + let mut has_err = 0; + let mut agree = 0; + + for (i, r) in readers.iter_mut().enumerate() { + if errs[i].is_some() { + has_err += 1; continue; } - }; - // If no current, add it. - if current.name.is_empty() { - top_entries.insert(i, entry.clone()); + + let entry = match r.peek().await { + Ok(res) => { + if let Some(entry) = res { + entry + } else { + // eof + at_eof += 1; + + continue; + } + } + Err(err) => { + if is_err_eof(&err) { + at_eof += 1; + continue; + } else if is_err_file_not_found(&err) { + at_eof += 1; + fnf += 1; + continue; + } else if is_err_volume_not_found(&err) { + at_eof += 1; + fnf += 1; + vnf += 1; + continue; + } else { + has_err += 1; + errs[i] = Some(err); + continue; + } + } + }; + + // If no current, add it. + if current.name.is_empty() { + top_entries.insert(i, Some(entry.clone())); + current = entry; + agree += 1; + + continue; + } + // If exact match, we agree. + if let Ok((_, true)) = current.matches(&entry, true) { + top_entries.insert(i, Some(entry)); + agree += 1; + + continue; + } + // If only the name matches we didn't agree, but add it for resolution. + if entry.name == current.name { + top_entries.insert(i, Some(entry)); + + continue; + } + // We got different entries + if entry.name > current.name { + continue; + } + // We got a new, better current. + // Clear existing entries. + top_entries.clear(); + agree += 1; + top_entries.insert(i, Some(entry.clone())); current = entry; - agree += 1; - continue; } - // If exact match, we agree. - if let Ok((_, true)) = current.matches(&entry, true) { - top_entries.insert(i, entry); - agree += 1; - continue; - } - // If only the name matches we didn't agree, but add it for resolution. - if entry.name == current.name { - top_entries.insert(i, entry); - continue; - } - // We got different entries - if entry.name > current.name { - continue; - } - // We got a new, better current. - // Clear existing entries. - top_entries.clear(); - agree += 1; - top_entries.insert(i, entry.clone()); - current = entry; - } - if has_err > 0 && has_err > opts.disks.len() - opts.min_disks { - if let Some(finished_fn) = opts.finished.as_ref() { - finished_fn(&errs).await; + if vnf > 0 && vnf >= (readers.len() - opts.min_disks) { + return Err(Error::new(DiskError::VolumeNotFound)); } - let mut combined_err = Vec::new(); - errs.iter().zip(opts.disks.iter()).for_each(|(err, disk)| match (err, disk) { - (Some(err), Some(disk)) => { - combined_err.push(format!("drive {} returned: {}", disk.to_string(), err)); + + if fnf > 0 && fnf >= (readers.len() - opts.min_disks) { + return Err(Error::new(DiskError::FileNotFound)); + } + + if has_err > 0 && has_err > opts.disks.len() - opts.min_disks { + if let Some(finished_fn) = opts.finished.as_ref() { + finished_fn(&errs).await; } - (Some(err), None) => { - combined_err.push(err.to_string()); + let mut combined_err = Vec::new(); + errs.iter().zip(opts.disks.iter()).for_each(|(err, disk)| match (err, disk) { + (Some(err), Some(disk)) => { + combined_err.push(format!("drive {} returned: {}", disk.to_string(), err)); + } + (Some(err), None) => { + combined_err.push(err.to_string()); + } + _ => {} + }); + + return Err(Error::from_string(combined_err.join(", "))); + } + + // Break if all at EOF or error. + if at_eof + has_err == readers.len() { + if has_err > 0 { + if let Some(finished_fn) = opts.finished.as_ref() { + if has_err > 0 { + finished_fn(&errs).await; + } + } } - _ => {} - }); - return Err(Error::from_string(combined_err.join(", "))); - } + break; + } - // Break if all at EOF or error. - if at_eof + has_err == readers.len() && has_err > 0 { - if let Some(finished_fn) = opts.finished.as_ref() { - finished_fn(&errs).await; + if agree == readers.len() { + for r in readers.iter_mut() { + let _ = r.skip(1).await; + } + if let Some(agreed_fn) = opts.agreed.as_ref() { + agreed_fn(current).await; + } + + continue; + } + + for (i, r) in readers.iter_mut().enumerate() { + if top_entries[i].is_some() { + let _ = r.skip(1).await; + } + } + + if let Some(partial_fn) = opts.partial.as_ref() { + partial_fn(MetaCacheEntries(top_entries), &errs).await; } break; } + Ok(()) + }); - if agree == readers.len() { - if let Some(agreed_fn) = opts.agreed.as_ref() { - agreed_fn(current).await; - } - continue; - } + jobs.push(revjob); - if let Some(partial_fn) = opts.partial.as_ref() { - partial_fn(MetaCacheEntries(top_entries), &errs).await; - } - } + let a = join_all(jobs).await; Ok(()) } diff --git a/ecstore/src/disk/error.rs b/ecstore/src/disk/error.rs index d20c00a3f..244c5928f 100644 --- a/ecstore/src/disk/error.rs +++ b/ecstore/src/disk/error.rs @@ -273,6 +273,17 @@ pub fn is_err_file_not_found(err: &Error) -> bool { matches!(err.downcast_ref::(), Some(DiskError::FileNotFound)) } +pub fn is_err_volume_not_found(err: &Error) -> bool { + matches!(err.downcast_ref::(), Some(DiskError::VolumeNotFound)) +} + +pub fn is_err_eof(err: &Error) -> bool { + if let Some(ioerr) = err.downcast_ref::() { + return ioerr.kind() == ErrorKind::UnexpectedEof; + } + false +} + pub fn is_sys_err_no_space(e: &io::Error) -> bool { if let Some(no) = e.raw_os_error() { return no == 28; diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index fa82685af..131630b3f 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -2426,51 +2426,4 @@ mod test { let _ = fs::remove_dir_all(&p).await; } - - #[tokio::test] - async fn test_walk_dir() { - let mut ep = Endpoint::try_from("/Users/weisd/project/weisd/s3-rustfs/target/volume/test").unwrap(); - ep.pool_idx = 0; - ep.set_idx = 0; - ep.disk_idx = 0; - - let disk = match LocalDisk::new(&ep, false).await { - Ok(res) => res, - Err(err) => { - println!("LocalDisk::new err {:?}", err); - return; - } - }; - - let (rd, mut wr) = tokio::io::duplex(1024); - - let job = tokio::spawn(async move { - let opts = WalkDirOptions { - bucket: "dada".to_owned(), - base_dir: "".to_owned(), - recursive: true, - ..Default::default() - }; - if let Err(err) = disk.walk_dir(opts, &mut wr).await { - println!("walk_dir err {:?}", err); - } - }); - - let job2 = tokio::spawn(async move { - let mut mrd = MetacacheReader::new(rd); - - while let Some(info) = mrd.peek().await.unwrap_or_default() { - println!("info {:?}", info.name) - } - }); - job.await; - job2.await; - - // let mut jobs = Vec::new(); - - // jobs.push(job); - // jobs.push(job2); - - // join_all(jobs).await; - } } diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index c47dd9afe..b1cc35d9b 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -770,7 +770,8 @@ impl MetaCacheEntry { } } -pub struct MetaCacheEntries(pub Vec); +#[derive(Debug)] +pub struct MetaCacheEntries(pub Vec>); impl MetaCacheEntries { pub fn resolve(&self, mut params: MetadataResolutionParams) -> Result> { @@ -785,7 +786,7 @@ impl MetaCacheEntries { let mut objs_agree = 0; let mut objs_valid = 0; - for entry in self.0.iter() { + for entry in self.0.iter().flatten() { if entry.name.is_empty() { continue; } @@ -859,7 +860,7 @@ impl MetaCacheEntries { } pub fn first_found(&self) -> (Option, usize) { - (self.0.iter().find(|x| !x.name.is_empty()).cloned(), self.0.len()) + (self.0.iter().find(|x| x.is_some()).cloned().unwrap_or_default(), self.0.len()) } } diff --git a/ecstore/src/endpoints.rs b/ecstore/src/endpoints.rs index b1205b87d..334c74d65 100644 --- a/ecstore/src/endpoints.rs +++ b/ecstore/src/endpoints.rs @@ -1,5 +1,4 @@ -use tracing::{info, warn}; -use url::Url; +use tracing::warn; use crate::{ disk::endpoint::{Endpoint, EndpointType}, diff --git a/ecstore/src/error.rs b/ecstore/src/error.rs index e7912cd8e..3d32b495c 100644 --- a/ecstore/src/error.rs +++ b/ecstore/src/error.rs @@ -1,9 +1,6 @@ -use std::io; - -use tracing::warn; -use tracing_error::{SpanTrace, SpanTraceStatus}; - use crate::disk::error::{clone_disk_err, DiskError}; +use std::io; +use tracing_error::{SpanTrace, SpanTraceStatus}; pub type StdError = Box; diff --git a/ecstore/src/file_meta.rs b/ecstore/src/file_meta.rs index 3dae801d6..8f6b03ea6 100644 --- a/ecstore/src/file_meta.rs +++ b/ecstore/src/file_meta.rs @@ -2096,10 +2096,6 @@ pub async fn read_xl_meta_no_data(reader: &mut R, size: us #[cfg(test)] mod test { - use std::fs; - use std::fs::File; - use std::os::unix::fs::MetadataExt; - use super::*; #[test] diff --git a/ecstore/src/metacache/writer.rs b/ecstore/src/metacache/writer.rs index c9e2ce599..c1bc1d98f 100644 --- a/ecstore/src/metacache/writer.rs +++ b/ecstore/src/metacache/writer.rs @@ -123,6 +123,8 @@ pub struct MetacacheReader { err: Option, buf: Vec, offset: usize, + + current: Option, } impl MetacacheReader { @@ -133,6 +135,7 @@ impl MetacacheReader { err: None, buf: Vec::new(), offset: 0, + current: None, } } @@ -199,7 +202,7 @@ impl MetacacheReader { 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?), - marker => Err(Error::msg("str marker err")), + _marker => Err(Error::msg("str marker err")), } } @@ -239,6 +242,45 @@ impl MetacacheReader { Ok(u32::from_be_bytes(buf.try_into().expect("Slice with incorrect length"))) } + pub async fn skip(&mut self, size: usize) -> Result<()> { + self.check_init().await?; + + if let Some(err) = &self.err { + return Err(err.clone()); + } + + let mut n = size; + + if self.current.is_some() { + n -= 1; + self.current = None; + } + + while n > 0 { + match rmp::decode::read_bool(&mut self.read_more(1).await?) { + Ok(res) => { + if !res { + return Ok(()); + } + } + Err(err) => { + let serr = format!("{:?}", err); + self.err = Some(Error::msg(&serr)); + return Err(Error::msg(&serr)); + } + }; + + let l = self.read_str_len().await?; + let _ = self.read_more(l as usize).await?; + let l = self.read_bin_len().await?; + let _ = self.read_more(l as usize).await?; + + n -= 1; + } + + Ok(()) + } + pub async fn peek(&mut self) -> Result> { self.check_init().await?; @@ -279,12 +321,15 @@ impl MetacacheReader { self.reset(); - Ok(Some(MetaCacheEntry { + let entry = Some(MetaCacheEntry { name, metadata, cached: None, reusable: false, - })) + }); + self.current = entry.clone(); + + Ok(entry) } pub async fn read_all(&mut self) -> Result> { diff --git a/ecstore/src/store_list_objects.rs b/ecstore/src/store_list_objects.rs index 5ff4586f1..fc7d47fe2 100644 --- a/ecstore/src/store_list_objects.rs +++ b/ecstore/src/store_list_objects.rs @@ -290,3 +290,132 @@ impl ECStore { Ok(ress) } } + +// list_path_raw + +#[cfg(test)] +mod test { + use crate::cache_value::metacache_set::list_path_raw; + use crate::cache_value::metacache_set::ListPathRawOptions; + use crate::disk::endpoint::Endpoint; + use crate::disk::error::is_err_eof; + use crate::disk::new_disk; + use crate::disk::DiskAPI; + use crate::disk::DiskOption; + use crate::disk::MetaCacheEntries; + use crate::disk::MetaCacheEntry; + use crate::disk::WalkDirOptions; + use crate::error::Error; + use crate::metacache::writer::MetacacheReader; + use futures::future::join_all; + use tokio::sync::broadcast; + + #[tokio::test] + async fn test_walk_dir() { + let mut ep = Endpoint::try_from("/Users/weisd/project/weisd/s3-rustfs/target/volume/test").unwrap(); + ep.pool_idx = 0; + ep.set_idx = 0; + ep.disk_idx = 0; + ep.is_local = true; + + let disk = new_disk(&ep, &DiskOption::default()).await.expect("init disk fail"); + + // let disk = match LocalDisk::new(&ep, false).await { + // Ok(res) => res, + // Err(err) => { + // println!("LocalDisk::new err {:?}", err); + // return; + // } + // }; + + let (rd, mut wr) = tokio::io::duplex(64); + + let job = tokio::spawn(async move { + let opts = WalkDirOptions { + bucket: "dada".to_owned(), + base_dir: "".to_owned(), + recursive: true, + ..Default::default() + }; + + println!("walk opts {:?}", opts); + if let Err(err) = disk.walk_dir(opts, &mut wr).await { + println!("walk_dir err {:?}", err); + } + }); + + let job2 = tokio::spawn(async move { + let mut mrd = MetacacheReader::new(rd); + + loop { + match mrd.peek().await { + Ok(res) => { + if let Some(info) = res { + println!("info {:?}", info.name) + } else { + break; + } + } + Err(err) => { + if is_err_eof(&err) { + break; + } + + println!("get err {:?}", err); + break; + } + } + } + }); + join_all(vec![job, job2]).await; + } + + #[tokio::test] + async fn test_list_path_raw() { + let mut ep = Endpoint::try_from("/Users/weisd/project/weisd/s3-rustfs/target/volume/test").unwrap(); + ep.pool_idx = 0; + ep.set_idx = 0; + ep.disk_idx = 0; + ep.is_local = true; + + let disk = new_disk(&ep, &DiskOption::default()).await.expect("init disk fail"); + + // let disk = match LocalDisk::new(&ep, false).await { + // Ok(res) => res, + // Err(err) => { + // println!("LocalDisk::new err {:?}", err); + // return; + // } + // }; + + let (_, rx) = broadcast::channel(1); + let bucket = "dada".to_owned(); + let forward_to = "".to_owned(); + let disks = vec![Some(disk)]; + let fallback_disks = Vec::new(); + + list_path_raw( + rx, + ListPathRawOptions { + disks, + fallback_disks, + bucket, + path: "".to_owned(), + recursice: true, + forward_to, + min_disks: 1, + report_not_found: false, + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + Box::pin(async move { println!("get entry: {}", entry.name) }) + })), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + Box::pin(async move { println!("get entries: {:?}", entries) }) + })), + finished: None, + ..Default::default() + }, + ) + .await + .unwrap(); + } +}