diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 5a8979bb4..2fe011da9 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -37,7 +37,10 @@ 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::{self, clean, decode_dir_object, has_suffix, path_join, GLOBAL_DIR_SUFFIX_WITH_SLASH, SLASH_SEPARATOR}; +use crate::utils::path::{ + self, clean, decode_dir_object, has_suffix, path_join, path_join_buf, GLOBAL_DIR_SUFFIX, GLOBAL_DIR_SUFFIX_WITH_SLASH, + SLASH_SEPARATOR, +}; use crate::{ file_meta::FileMeta, store_api::{FileInfo, RawFileInfo}, @@ -743,7 +746,9 @@ impl LocalDisk { let mut entries = match self.list_dir("", &opts.bucket, current, -1).await { Ok(res) => res, Err(e) => { - if !DiskError::VolumeNotFound.is(&e) && !is_err_file_not_found(&e) {} + if !DiskError::VolumeNotFound.is(&e) && !is_err_file_not_found(&e) { + error!("scan list_dir {}, err {:?}", ¤t, &e); + } if opts.report_notfound && is_err_file_not_found(&e) && current == &opts.base_dir { return Err(Error::new(DiskError::FileNotFound)); @@ -849,13 +854,13 @@ impl LocalDisk { if pop < name { // out.write_obj(&MetaCacheEntry { - name: pop, + name: pop.clone(), ..Default::default() }) .await?; if opts.recursive { - if let Err(er) = Box::pin(self.scan_dir(current, opts, out, objs_returned)).await { + if let Err(er) = Box::pin(self.scan_dir(&mut pop.clone(), opts, out, objs_returned)).await { error!("scan_dir err {:?}", er); } } @@ -1495,7 +1500,11 @@ impl DiskAPI for LocalDisk { } let volume_dir = self.get_bucket_path(volume)?; - let dir_path_abs = volume_dir.join(Path::new(&dir_path)); + let dir_path_abs = volume_dir.join(Path::new(&dir_path.trim_start_matches(SLASH_SEPARATOR))); + println!( + "list dir volume_dir: {:?} join dir_path {} = abs {:?}", + &volume_dir, &dir_path, dir_path_abs + ); let entries = match os::read_dir(&dir_path_abs, count).await { Ok(res) => res, @@ -1534,8 +1543,14 @@ impl DiskAPI for LocalDisk { if opts.base_dir.ends_with(SLASH_SEPARATOR) { let fpath = self.get_object_path( &opts.bucket, - format!("{}/{}", opts.base_dir.trim_end_matches(SLASH_SEPARATOR), STORAGE_FORMAT_FILE).as_str(), + path_join_buf(&[ + format!("{}{}", opts.base_dir.trim_end_matches(SLASH_SEPARATOR), GLOBAL_DIR_SUFFIX).as_str(), + STORAGE_FORMAT_FILE, + ]) + .as_str(), )?; + + println!("fpath {:?}", &fpath); if let Ok(data) = self.read_metadata(fpath).await { let meta = MetaCacheEntry { name: opts.base_dir.clone(), @@ -2434,8 +2449,6 @@ mod test { 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(), @@ -2447,8 +2460,6 @@ mod test { } }); - // let rd = VecAsyncReader::new(wr.get_buffer().to_vec()); - let rd_job = tokio::spawn(async move { let mut mrd = MetacacheReader::new(rd); diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index a6b182ce8..e1ed0b48d 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -163,7 +163,7 @@ impl VecAsyncWriter { // Implementing AsyncWrite trait for VecAsyncWriter impl AsyncWrite for VecAsyncWriter { - fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll> { + fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll> { let len = buf.len(); // Assume synchronous writing for simplicity @@ -178,7 +178,7 @@ impl AsyncWrite for VecAsyncWriter { Poll::Ready(Ok(())) } - fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { // Similar to flush, shutdown has no effect here Poll::Ready(Ok(())) } diff --git a/ecstore/src/metacache/writer.rs b/ecstore/src/metacache/writer.rs index ddc55f913..a405722e5 100644 --- a/ecstore/src/metacache/writer.rs +++ b/ecstore/src/metacache/writer.rs @@ -71,6 +71,8 @@ impl MetacacheWriter { } pub async fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> { + println!("write_obj {:?}", &obj); + self.init().await?; rmp::encode::write_bool(&mut self.buf, true).map_err(|e| Error::msg(format!("{:?}", e)))?; @@ -162,8 +164,6 @@ impl MetacacheReader { let data = &self.buf[pref..ext_size]; - println!("pref {} offset {},ext_size {}, data {:?}", pref, self.offset, ext_size, &data); - Ok(data) } @@ -181,8 +181,6 @@ impl MetacacheReader { 0 } }; - - println!("ver {}", ver); match ver { 1 | 2 => (), _ => { @@ -233,21 +231,25 @@ impl MetacacheReader { } async fn read_u8(&mut self) -> Result { - let a = self.read_more(1).await?; + let buf = self.read_more(1).await?; - Ok(a[0]) + Ok(u8::from_be_bytes(buf.try_into().expect("Slice with incorrect length"))) } async fn read_u16(&mut self) -> Result { - rmp::decode::read_u16(&mut self.read_more(2).await?).map_err(|e| Error::msg(format!("{:?}", e))) + let buf = self.read_more(2).await?; + + Ok(u16::from_be_bytes(buf.try_into().expect("Slice with incorrect length"))) } async fn read_u32(&mut self) -> Result { - rmp::decode::read_u32(&mut self.read_more(4).await?).map_err(|e| Error::msg(format!("{:?}", e))) + let buf = self.read_more(4).await?; + + Ok(u32::from_be_bytes(buf.try_into().expect("Slice with incorrect length"))) } pub async fn peek(&mut self) -> Result> { - self.check_init().await; + self.check_init().await?; if let Some(err) = &self.err { return Err(err.clone()); @@ -337,8 +339,6 @@ async fn test_writer() { let data = f.get_buffer().to_vec(); - println!("data len {}", data.len()); - let nf = VecAsyncReader::new(data); let mut r = MetacacheReader::new(nf); diff --git a/ecstore/src/utils/path.rs b/ecstore/src/utils/path.rs index 9f2c22a08..745aef679 100644 --- a/ecstore/src/utils/path.rs +++ b/ecstore/src/utils/path.rs @@ -1,7 +1,7 @@ use std::path::Path; use std::path::PathBuf; -const GLOBAL_DIR_SUFFIX: &str = "__XLDIR__"; +pub const GLOBAL_DIR_SUFFIX: &str = "__XLDIR__"; pub const SLASH_SEPARATOR: &str = "/";