From 5501649d7974781a1e452570344d633c00a0f6eb Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 9 Dec 2024 10:03:02 +0800 Subject: [PATCH] scandir --- ecstore/src/disk/local.rs | 255 +++++++++++++++++++++----------- ecstore/src/metacache/writer.rs | 74 +++++++-- 2 files changed, 236 insertions(+), 93 deletions(-) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 5994cfa3c..8eec78e08 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -26,6 +26,8 @@ use crate::heal::data_usage_cache::{DataUsageCache, DataUsageEntry}; use crate::heal::error::{ERR_IGNORE_FILE_CONTRIB, ERR_SKIP_FILE}; use crate::heal::heal_commands::{HealScanMode, HealingTracker}; use crate::heal::heal_ops::HEALING_TRACKER_FILENAME; +use crate::io::AsyncToSync; +use crate::metacache::writer::MetacacheWriter; use crate::new_object_layer_fn; use crate::set_disk::{ conv_part_err_to_int, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN, @@ -34,7 +36,7 @@ use crate::set_disk::{ use crate::store_api::{BitrotAlgorithm, StorageAPI}; use crate::utils::fs::{access, lstat, O_APPEND, O_CREATE, O_RDONLY, O_WRONLY}; use crate::utils::os::get_info; -use crate::utils::path::{clean, has_suffix, path_join, GLOBAL_DIR_SUFFIX_WITH_SLASH, SLASH_SEPARATOR}; +use crate::utils::path::{clean, decode_dir_object, has_suffix, path_join, GLOBAL_DIR_SUFFIX_WITH_SLASH, SLASH_SEPARATOR}; use crate::{ file_meta::FileMeta, store_api::{FileInfo, RawFileInfo}, @@ -45,7 +47,7 @@ use nix::NixPath; use path_absolutize::Absolutize; use std::collections::{HashMap, HashSet}; use std::fmt::Debug; -use std::io::Cursor; +use std::io::{Cursor, Write}; use std::os::unix::fs::MetadataExt; use std::sync::atomic::{AtomicU32, Ordering}; use std::sync::Arc; @@ -59,7 +61,7 @@ use tokio::fs::{self, File}; use tokio::io::{AsyncReadExt, 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)] @@ -712,6 +714,162 @@ impl LocalDisk { let n = file.read_to_end(&mut data).await?; bitrot_verify(&mut Cursor::new(data), n, part_size, algo, sum.to_vec(), shard_size) } + + async fn scan_dir( + &self, + current: &mut String, + opts: &WalkDirOptions, + out: &mut MetacacheWriter, + objs_returned: &mut i32, + ) -> Result<()> { + let forward = { + if !opts.forward_to.is_empty() && opts.forward_to.starts_with(&*current) { + let forward = opts.forward_to.trim_start_matches(&*current); + if let Some(idx) = forward.find('/') { + &forward[..idx] + } else { + forward + } + } else { + "" + } + }; + + if opts.limit > 0 && *objs_returned >= opts.limit { + return Ok(()); + } + + let mut entries = match self.list_dir("", &opts.bucket, &opts.base_dir, -1).await { + Ok(res) => res, + Err(e) => { + if !DiskError::VolumeNotFound.is(&e) && !is_err_file_not_found(&e) { + info!("list_dir err {:?}", &e); + } + + if opts.report_notfound && is_err_file_not_found(&e) { + return Err(e); + } + return Ok(()); + } + }; + + if entries.is_empty() { + return Ok(()); + } + + let s = SLASH_SEPARATOR.chars().next().unwrap_or_default(); + *current = current.trim_matches(s).to_owned(); + + let bucket = opts.bucket.as_str(); + + let mut dir_objes = HashSet::new(); + + // 第一层过滤 + for item in entries.iter_mut() { + // warn!("walk_dir get entry {:?}", &entry); + + let entry = item.clone(); + // check limit + if opts.limit > 0 && *objs_returned >= opts.limit { + return Ok(()); + } + // check prefix + if !opts.filter_prefix.is_empty() && !entry.starts_with(&opts.filter_prefix) { + *item = "".to_owned(); + continue; + } + + if !forward.is_empty() && entry.as_str() < forward { + *item = "".to_owned(); + continue; + } + + if entry.ends_with(SLASH_SEPARATOR) { + if entry.ends_with(GLOBAL_DIR_SUFFIX_WITH_SLASH) { + let entry = format!("{}{}", entry.as_str().trim_end_matches(GLOBAL_DIR_SUFFIX_WITH_SLASH), SLASH_SEPARATOR); + dir_objes.insert(entry.clone()); + *item = entry; + continue; + } + + *item = entry.trim_end_matches(SLASH_SEPARATOR).to_owned(); + continue; + } + + *item = "".to_owned(); + + if entry.ends_with(STORAGE_FORMAT_FILE) { + // + let metadata = self + .read_metadata(self.get_object_path(bucket, format!("{}/{}", ¤t, &entry).as_str())?) + .await?; + let name = entry.trim_end_matches(STORAGE_FORMAT_FILE).trim_end_matches(SLASH_SEPARATOR); + let name = decode_dir_object(format!("{}/{}", ¤t, &name).as_str()); + + out.write_obj(&MetaCacheEntry { + name, + metadata, + ..Default::default() + })?; + *objs_returned += 1; + return Ok(()); + } + } + + entries.sort(); + + let mut entries = entries.as_slice(); + if !forward.is_empty() { + for (i, entry) in entries.iter().enumerate() { + if entry.as_str() >= forward || forward.starts_with(entry.as_str()) { + entries = &entries[i..]; + break; + } + } + } + + let prefix = ""; + let mut dir_stack: Vec = Vec::with_capacity(5); + + for entry in entries.iter() { + // + if opts.limit > 0 && *objs_returned >= opts.limit { + return Ok(()); + } + + if entry.is_empty() { + continue; + } + + let name = format!("{}{}", current, entry); + + if !dir_stack.is_empty() { + if let Some(pop) = dir_stack.pop() { + if pop < name { + // + out.write_obj(&MetaCacheEntry { + name: pop, + ..Default::default() + })?; + + if opts.recursive { + if let Err(er) = Box::pin(self.scan_dir(current, opts, out, objs_returned)).await { + error!("scan_dir err {:?}", er); + } + } + } + } + } + + if let Some(dir) = dir_objes.get(entry) { + // + } + } + + // FIXME:TODO: + + Ok(()) + } } fn is_root_path(path: impl AsRef) -> bool { @@ -1293,7 +1451,7 @@ impl DiskAPI for LocalDisk { Ok(entries) } - // FIXME: TODO: io.writer + // FIXME: TODO: io.writer TODO cancel async fn walk_dir(&self, opts: WalkDirOptions, wr: crate::io::Writer) -> Result> { // warn!("walk_dir opts {:?}", &opts); @@ -1305,7 +1463,11 @@ impl DiskAPI for LocalDisk { } } - let mut metas = Vec::new(); + let mut wr = wr; + + let mut out = MetacacheWriter::new(AsyncToSync::new_writer(&mut wr)); + + let mut objs_returned = 0; if opts.base_dir.ends_with(SLASH_SEPARATOR) { let fpath = self.get_object_path( @@ -1318,87 +1480,14 @@ impl DiskAPI for LocalDisk { metadata: data, ..Default::default() }; - metas.push(meta); - return Ok(metas); + out.write_obj(&meta)?; + objs_returned += 1; } } - let mut entries = match self.list_dir("", &opts.bucket, &opts.base_dir, -1).await { - Ok(res) => res, - Err(e) => { - if !DiskError::VolumeNotFound.is(&e) && !is_err_file_not_found(&e) { - info!("list_dir err {:?}", &e); - } - - if opts.report_notfound && is_err_file_not_found(&e) { - return Err(e); - } - return Ok(Vec::new()); - } - }; - - if entries.is_empty() { - return Ok(Vec::new()); - } - - entries.sort(); - - // 已读计数 - let objs_returned = 0; - - let bucket = opts.bucket.as_str(); - - let mut dir_objes = HashSet::new(); - - // 第一层过滤 - for entry in entries.iter() { - // warn!("walk_dir get entry {:?}", &entry); - - // check limit - if opts.limit > 0 && objs_returned >= opts.limit { - return Ok(metas); - } - // check prefix - if !opts.filter_prefix.is_empty() && !entry.starts_with(&opts.filter_prefix) { - continue; - } - - let mut meta = MetaCacheEntry { ..Default::default() }; - - let mut name = { - if opts.base_dir.is_empty() { - entry.clone() - } else { - format!("{}{}{}", opts.base_dir.trim_end_matches(SLASH_SEPARATOR), SLASH_SEPARATOR, entry) - } - }; - - if name.ends_with(SLASH_SEPARATOR) { - if name.ends_with(GLOBAL_DIR_SUFFIX_WITH_SLASH) { - name = format!("{}{}", name.as_str().trim_end_matches(GLOBAL_DIR_SUFFIX_WITH_SLASH), SLASH_SEPARATOR); - dir_objes.insert(name.clone()); - } else { - name = name.as_str().trim_end_matches(SLASH_SEPARATOR).to_owned(); - } - } - meta.name = name; - - let fpath = self.get_object_path(bucket, format!("{}/{}", &meta.name, STORAGE_FORMAT_FILE).as_str())?; - - if let Ok(data) = self.read_metadata(&fpath).await { - meta.metadata = data; - } else { - let fpath = self.get_object_path(bucket, &meta.name)?; - - if !is_empty_dir(fpath).await { - meta.name = format!("{}{}", &meta.name, SLASH_SEPARATOR); - } - } - - metas.push(meta); - } - - Ok(metas) + let mut current = opts.base_dir.clone(); + self.scan_dir(&mut current, &opts, &mut out, &mut objs_returned).await?; + Ok(Vec::new()) } // #[tracing::instrument(skip(self))] diff --git a/ecstore/src/metacache/writer.rs b/ecstore/src/metacache/writer.rs index 57b0d189e..cc23cf6ec 100644 --- a/ecstore/src/metacache/writer.rs +++ b/ecstore/src/metacache/writer.rs @@ -4,17 +4,35 @@ use crate::error::Result; use std::io::Read; use std::io::Write; use std::str::from_utf8; +// use std::sync::Arc; +// use tokio::sync::mpsc; +// use tokio::sync::mpsc::Sender; +// use tokio::task; const METACACHE_STREAM_VERSION: u8 = 2; +#[derive(Debug)] pub struct MetacacheWriter { wr: W, created: bool, + // err: Option, } impl MetacacheWriter { pub fn new(wr: W) -> Self { - Self { wr, created: false } + Self { + wr, + created: false, + // err: None, + } + } + + pub fn init(&mut self) -> Result<()> { + if !self.created { + rmp::encode::write_u8(&mut self.wr, METACACHE_STREAM_VERSION)?; + self.created = true; + } + Ok(()) } pub fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> { @@ -22,26 +40,62 @@ impl MetacacheWriter { return Ok(()); } - if !self.created { - rmp::encode::write_u8(&mut self.wr, METACACHE_STREAM_VERSION)?; - self.created = false; - } + self.init()?; 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)?; + self.write_obj(obj)?; } Ok(()) } + pub fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> { + self.init()?; + 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(()) + } + + // pub async fn stream(&mut self) -> Result> { + // let (sender, mut receiver) = mpsc::channel::(100); + + // let wr = Arc::new(self); + + // task::spawn(async move { + // while let Some(obj) = receiver.recv().await { + // // if obj.name.is_empty() || self.err.is_some() { + // // continue; + // // } + + // let _ = wr.write_obj(&obj); + + // // if let Err(err) = rmp::encode::write_bool(&mut self.wr, true) { + // // self.err = Some(Error::new(err)); + // // continue; + // // } + + // // if let Err(err) = rmp::encode::write_str(&mut self.wr, &obj.name) { + // // self.err = Some(Error::new(err)); + // // continue; + // // } + + // // if let Err(err) = rmp::encode::write_bin(&mut self.wr, &obj.metadata) { + // // self.err = Some(Error::new(err)); + // // continue; + // // } + // } + // }); + + // Ok(sender) + // } + pub fn close(&mut self) -> Result<()> { rmp::encode::write_bool(&mut self.wr, false)?; self.wr.flush()?;