mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
test walk_dir done
This commit is contained in:
+21
-10
@@ -37,7 +37,10 @@ use crate::set_disk::{
|
|||||||
use crate::store_api::{BitrotAlgorithm, StorageAPI};
|
use crate::store_api::{BitrotAlgorithm, StorageAPI};
|
||||||
use crate::utils::fs::{access, lstat, O_APPEND, O_CREATE, O_RDONLY, O_WRONLY};
|
use crate::utils::fs::{access, lstat, O_APPEND, O_CREATE, O_RDONLY, O_WRONLY};
|
||||||
use crate::utils::os::get_info;
|
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::{
|
use crate::{
|
||||||
file_meta::FileMeta,
|
file_meta::FileMeta,
|
||||||
store_api::{FileInfo, RawFileInfo},
|
store_api::{FileInfo, RawFileInfo},
|
||||||
@@ -743,7 +746,9 @@ impl LocalDisk {
|
|||||||
let mut entries = match self.list_dir("", &opts.bucket, current, -1).await {
|
let mut entries = match self.list_dir("", &opts.bucket, current, -1).await {
|
||||||
Ok(res) => res,
|
Ok(res) => res,
|
||||||
Err(e) => {
|
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 {
|
if opts.report_notfound && is_err_file_not_found(&e) && current == &opts.base_dir {
|
||||||
return Err(Error::new(DiskError::FileNotFound));
|
return Err(Error::new(DiskError::FileNotFound));
|
||||||
@@ -849,13 +854,13 @@ impl LocalDisk {
|
|||||||
if pop < name {
|
if pop < name {
|
||||||
//
|
//
|
||||||
out.write_obj(&MetaCacheEntry {
|
out.write_obj(&MetaCacheEntry {
|
||||||
name: pop,
|
name: pop.clone(),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
})
|
})
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
if opts.recursive {
|
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);
|
error!("scan_dir err {:?}", er);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1495,7 +1500,11 @@ impl DiskAPI for LocalDisk {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let volume_dir = self.get_bucket_path(volume)?;
|
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 {
|
let entries = match os::read_dir(&dir_path_abs, count).await {
|
||||||
Ok(res) => res,
|
Ok(res) => res,
|
||||||
@@ -1534,8 +1543,14 @@ impl DiskAPI for LocalDisk {
|
|||||||
if opts.base_dir.ends_with(SLASH_SEPARATOR) {
|
if opts.base_dir.ends_with(SLASH_SEPARATOR) {
|
||||||
let fpath = self.get_object_path(
|
let fpath = self.get_object_path(
|
||||||
&opts.bucket,
|
&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 {
|
if let Ok(data) = self.read_metadata(fpath).await {
|
||||||
let meta = MetaCacheEntry {
|
let meta = MetaCacheEntry {
|
||||||
name: opts.base_dir.clone(),
|
name: opts.base_dir.clone(),
|
||||||
@@ -2434,8 +2449,6 @@ mod test {
|
|||||||
|
|
||||||
let (rd, mut wr) = tokio::io::duplex(64);
|
let (rd, mut wr) = tokio::io::duplex(64);
|
||||||
|
|
||||||
// let mut wr = VecAsyncWriter::new(Vec::new());
|
|
||||||
|
|
||||||
let job = tokio::spawn(async move {
|
let job = tokio::spawn(async move {
|
||||||
let opts = WalkDirOptions {
|
let opts = WalkDirOptions {
|
||||||
bucket: "dada".to_owned(),
|
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 rd_job = tokio::spawn(async move {
|
||||||
let mut mrd = MetacacheReader::new(rd);
|
let mut mrd = MetacacheReader::new(rd);
|
||||||
|
|
||||||
|
|||||||
+2
-2
@@ -163,7 +163,7 @@ impl VecAsyncWriter {
|
|||||||
|
|
||||||
// Implementing AsyncWrite trait for VecAsyncWriter
|
// Implementing AsyncWrite trait for VecAsyncWriter
|
||||||
impl AsyncWrite for VecAsyncWriter {
|
impl AsyncWrite for VecAsyncWriter {
|
||||||
fn poll_write(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll<io::Result<usize>> {
|
fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll<io::Result<usize>> {
|
||||||
let len = buf.len();
|
let len = buf.len();
|
||||||
|
|
||||||
// Assume synchronous writing for simplicity
|
// Assume synchronous writing for simplicity
|
||||||
@@ -178,7 +178,7 @@ impl AsyncWrite for VecAsyncWriter {
|
|||||||
Poll::Ready(Ok(()))
|
Poll::Ready(Ok(()))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
|
fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<io::Result<()>> {
|
||||||
// Similar to flush, shutdown has no effect here
|
// Similar to flush, shutdown has no effect here
|
||||||
Poll::Ready(Ok(()))
|
Poll::Ready(Ok(()))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -71,6 +71,8 @@ impl<W: AsyncWrite + Unpin> MetacacheWriter<W> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub async fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> {
|
pub async fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> {
|
||||||
|
println!("write_obj {:?}", &obj);
|
||||||
|
|
||||||
self.init().await?;
|
self.init().await?;
|
||||||
|
|
||||||
rmp::encode::write_bool(&mut self.buf, true).map_err(|e| Error::msg(format!("{:?}", e)))?;
|
rmp::encode::write_bool(&mut self.buf, true).map_err(|e| Error::msg(format!("{:?}", e)))?;
|
||||||
@@ -162,8 +164,6 @@ impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
|
|
||||||
let data = &self.buf[pref..ext_size];
|
let data = &self.buf[pref..ext_size];
|
||||||
|
|
||||||
println!("pref {} offset {},ext_size {}, data {:?}", pref, self.offset, ext_size, &data);
|
|
||||||
|
|
||||||
Ok(data)
|
Ok(data)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -181,8 +181,6 @@ impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
0
|
0
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
println!("ver {}", ver);
|
|
||||||
match ver {
|
match ver {
|
||||||
1 | 2 => (),
|
1 | 2 => (),
|
||||||
_ => {
|
_ => {
|
||||||
@@ -233,21 +231,25 @@ impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn read_u8(&mut self) -> Result<u8> {
|
async fn read_u8(&mut self) -> Result<u8> {
|
||||||
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<u16> {
|
async fn read_u16(&mut self) -> Result<u16> {
|
||||||
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<u32> {
|
async fn read_u32(&mut self) -> Result<u32> {
|
||||||
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<Option<MetaCacheEntry>> {
|
pub async fn peek(&mut self) -> Result<Option<MetaCacheEntry>> {
|
||||||
self.check_init().await;
|
self.check_init().await?;
|
||||||
|
|
||||||
if let Some(err) = &self.err {
|
if let Some(err) = &self.err {
|
||||||
return Err(err.clone());
|
return Err(err.clone());
|
||||||
@@ -337,8 +339,6 @@ async fn test_writer() {
|
|||||||
|
|
||||||
let data = f.get_buffer().to_vec();
|
let data = f.get_buffer().to_vec();
|
||||||
|
|
||||||
println!("data len {}", data.len());
|
|
||||||
|
|
||||||
let nf = VecAsyncReader::new(data);
|
let nf = VecAsyncReader::new(data);
|
||||||
|
|
||||||
let mut r = MetacacheReader::new(nf);
|
let mut r = MetacacheReader::new(nf);
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
|
|
||||||
const GLOBAL_DIR_SUFFIX: &str = "__XLDIR__";
|
pub const GLOBAL_DIR_SUFFIX: &str = "__XLDIR__";
|
||||||
|
|
||||||
pub const SLASH_SEPARATOR: &str = "/";
|
pub const SLASH_SEPARATOR: &str = "/";
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user