mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-04 03:05:39 +00:00
test_list_path_raw done
This commit is contained in:
@@ -1,5 +1,6 @@
|
|||||||
use std::{future::Future, pin::Pin, sync::Arc};
|
use std::{future::Future, pin::Pin, sync::Arc};
|
||||||
|
|
||||||
|
use futures::{future::join_all, join};
|
||||||
use tokio::{
|
use tokio::{
|
||||||
spawn,
|
spawn,
|
||||||
sync::{
|
sync::{
|
||||||
@@ -8,11 +9,16 @@ use tokio::{
|
|||||||
RwLock,
|
RwLock,
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
|
use tracing::error;
|
||||||
|
|
||||||
use crate::{
|
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},
|
error::{Error, Result},
|
||||||
io::Writer,
|
io::Writer,
|
||||||
|
metacache::writer::MetacacheReader,
|
||||||
};
|
};
|
||||||
|
|
||||||
type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
|
type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
|
||||||
@@ -62,44 +68,40 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
|
|||||||
return Err(Error::from_string("list_path_raw: 0 drives provided"));
|
return Err(Error::from_string("list_path_raw: 0 drives provided"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let mut jobs: Vec<tokio::task::JoinHandle<std::result::Result<(), Error>>> = Vec::new();
|
||||||
let mut readers = Vec::with_capacity(opts.disks.len());
|
let mut readers = Vec::with_capacity(opts.disks.len());
|
||||||
let fds = Arc::new(RwLock::new(opts.fallback_disks.clone()));
|
let fds = Arc::new(RwLock::new(opts.fallback_disks.clone()));
|
||||||
for disk in opts.disks.iter() {
|
for disk in opts.disks.iter() {
|
||||||
let disk = disk.clone();
|
let opdisk = disk.clone();
|
||||||
let opts_clone = opts.clone();
|
let opts_clone = opts.clone();
|
||||||
let fds_clone = fds.clone();
|
let fds_clone = fds.clone();
|
||||||
let (m_tx, m_rx) = mpsc::channel::<MetaCacheEntry>(100);
|
// let (m_tx, m_rx) = mpsc::channel::<MetaCacheEntry>(100);
|
||||||
readers.push(m_rx);
|
// readers.push(m_rx);
|
||||||
spawn(async move {
|
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;
|
let mut need_fallback = false;
|
||||||
if disk.is_none() {
|
if let Some(disk) = opdisk {
|
||||||
need_fallback = true;
|
match disk.walk_dir(wakl_opts, &mut wr).await {
|
||||||
} else {
|
Ok(_res) => {}
|
||||||
match disk
|
Err(err) => {
|
||||||
.as_ref()
|
error!("walk dir err {:?}", &err);
|
||||||
.unwrap()
|
need_fallback = true;
|
||||||
.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;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
Err(_) => need_fallback = true,
|
|
||||||
}
|
}
|
||||||
|
} else {
|
||||||
|
need_fallback = true;
|
||||||
}
|
}
|
||||||
|
|
||||||
while need_fallback {
|
while need_fallback {
|
||||||
@@ -108,133 +110,203 @@ pub async fn list_path_raw(mut rx: B_Receiver<bool>, opts: ListPathRawOptions) -
|
|||||||
if fds_w.is_empty() {
|
if fds_w.is_empty() {
|
||||||
break None;
|
break None;
|
||||||
}
|
}
|
||||||
let fd = fds_w.remove(0);
|
|
||||||
if fd.is_some() && fd.as_ref().unwrap().is_online().await {
|
if let Some(fd) = fds_w.remove(0) {
|
||||||
break fd;
|
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;
|
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<Option<Error>> = vec![None; readers.len()];
|
let revjob = spawn(async move {
|
||||||
loop {
|
let mut errs: Vec<Option<Error>> = vec![None; readers.len()];
|
||||||
let mut current = MetaCacheEntry::default();
|
loop {
|
||||||
let (mut at_eof, mut has_err, mut agree) = (0, 0, 0);
|
let mut current = MetaCacheEntry::default();
|
||||||
if rx.try_recv().is_ok() {
|
|
||||||
return Err(Error::from_string("canceled"));
|
|
||||||
}
|
|
||||||
let mut top_entries: Vec<MetaCacheEntry> = Vec::with_capacity(readers.len());
|
|
||||||
// top_entries.clear();
|
|
||||||
|
|
||||||
for (i, r) in readers.iter_mut().enumerate() {
|
if rx.try_recv().is_ok() {
|
||||||
if errs[i].is_none() {
|
return Err(Error::from_string("canceled"));
|
||||||
has_err += 1;
|
|
||||||
continue;
|
|
||||||
}
|
}
|
||||||
let entry = match r.recv().await {
|
let mut top_entries: Vec<Option<MetaCacheEntry>> = vec![None; readers.len()];
|
||||||
Some(entry) => entry,
|
|
||||||
None => {
|
let mut at_eof = 0;
|
||||||
at_eof += 1;
|
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;
|
continue;
|
||||||
}
|
}
|
||||||
};
|
|
||||||
// If no current, add it.
|
let entry = match r.peek().await {
|
||||||
if current.name.is_empty() {
|
Ok(res) => {
|
||||||
top_entries.insert(i, entry.clone());
|
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;
|
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 vnf > 0 && vnf >= (readers.len() - opts.min_disks) {
|
||||||
if let Some(finished_fn) = opts.finished.as_ref() {
|
return Err(Error::new(DiskError::VolumeNotFound));
|
||||||
finished_fn(&errs).await;
|
|
||||||
}
|
}
|
||||||
let mut combined_err = Vec::new();
|
|
||||||
errs.iter().zip(opts.disks.iter()).for_each(|(err, disk)| match (err, disk) {
|
if fnf > 0 && fnf >= (readers.len() - opts.min_disks) {
|
||||||
(Some(err), Some(disk)) => {
|
return Err(Error::new(DiskError::FileNotFound));
|
||||||
combined_err.push(format!("drive {} returned: {}", disk.to_string(), err));
|
}
|
||||||
|
|
||||||
|
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) => {
|
let mut combined_err = Vec::new();
|
||||||
combined_err.push(err.to_string());
|
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 agree == readers.len() {
|
||||||
if at_eof + has_err == readers.len() && has_err > 0 {
|
for r in readers.iter_mut() {
|
||||||
if let Some(finished_fn) = opts.finished.as_ref() {
|
let _ = r.skip(1).await;
|
||||||
finished_fn(&errs).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;
|
break;
|
||||||
}
|
}
|
||||||
|
Ok(())
|
||||||
|
});
|
||||||
|
|
||||||
if agree == readers.len() {
|
jobs.push(revjob);
|
||||||
if let Some(agreed_fn) = opts.agreed.as_ref() {
|
|
||||||
agreed_fn(current).await;
|
|
||||||
}
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
if let Some(partial_fn) = opts.partial.as_ref() {
|
let a = join_all(jobs).await;
|
||||||
partial_fn(MetaCacheEntries(top_entries), &errs).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -273,6 +273,17 @@ pub fn is_err_file_not_found(err: &Error) -> bool {
|
|||||||
matches!(err.downcast_ref::<DiskError>(), Some(DiskError::FileNotFound))
|
matches!(err.downcast_ref::<DiskError>(), Some(DiskError::FileNotFound))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn is_err_volume_not_found(err: &Error) -> bool {
|
||||||
|
matches!(err.downcast_ref::<DiskError>(), Some(DiskError::VolumeNotFound))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn is_err_eof(err: &Error) -> bool {
|
||||||
|
if let Some(ioerr) = err.downcast_ref::<io::Error>() {
|
||||||
|
return ioerr.kind() == ErrorKind::UnexpectedEof;
|
||||||
|
}
|
||||||
|
false
|
||||||
|
}
|
||||||
|
|
||||||
pub fn is_sys_err_no_space(e: &io::Error) -> bool {
|
pub fn is_sys_err_no_space(e: &io::Error) -> bool {
|
||||||
if let Some(no) = e.raw_os_error() {
|
if let Some(no) = e.raw_os_error() {
|
||||||
return no == 28;
|
return no == 28;
|
||||||
|
|||||||
@@ -2426,51 +2426,4 @@ mod test {
|
|||||||
|
|
||||||
let _ = fs::remove_dir_all(&p).await;
|
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;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -770,7 +770,8 @@ impl MetaCacheEntry {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct MetaCacheEntries(pub Vec<MetaCacheEntry>);
|
#[derive(Debug)]
|
||||||
|
pub struct MetaCacheEntries(pub Vec<Option<MetaCacheEntry>>);
|
||||||
|
|
||||||
impl MetaCacheEntries {
|
impl MetaCacheEntries {
|
||||||
pub fn resolve(&self, mut params: MetadataResolutionParams) -> Result<Option<MetaCacheEntry>> {
|
pub fn resolve(&self, mut params: MetadataResolutionParams) -> Result<Option<MetaCacheEntry>> {
|
||||||
@@ -785,7 +786,7 @@ impl MetaCacheEntries {
|
|||||||
let mut objs_agree = 0;
|
let mut objs_agree = 0;
|
||||||
let mut objs_valid = 0;
|
let mut objs_valid = 0;
|
||||||
|
|
||||||
for entry in self.0.iter() {
|
for entry in self.0.iter().flatten() {
|
||||||
if entry.name.is_empty() {
|
if entry.name.is_empty() {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
@@ -859,7 +860,7 @@ impl MetaCacheEntries {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn first_found(&self) -> (Option<MetaCacheEntry>, usize) {
|
pub fn first_found(&self) -> (Option<MetaCacheEntry>, 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())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,5 +1,4 @@
|
|||||||
use tracing::{info, warn};
|
use tracing::warn;
|
||||||
use url::Url;
|
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
disk::endpoint::{Endpoint, EndpointType},
|
disk::endpoint::{Endpoint, EndpointType},
|
||||||
|
|||||||
@@ -1,9 +1,6 @@
|
|||||||
use std::io;
|
|
||||||
|
|
||||||
use tracing::warn;
|
|
||||||
use tracing_error::{SpanTrace, SpanTraceStatus};
|
|
||||||
|
|
||||||
use crate::disk::error::{clone_disk_err, DiskError};
|
use crate::disk::error::{clone_disk_err, DiskError};
|
||||||
|
use std::io;
|
||||||
|
use tracing_error::{SpanTrace, SpanTraceStatus};
|
||||||
|
|
||||||
pub type StdError = Box<dyn std::error::Error + Send + Sync + 'static>;
|
pub type StdError = Box<dyn std::error::Error + Send + Sync + 'static>;
|
||||||
|
|
||||||
|
|||||||
@@ -2096,10 +2096,6 @@ pub async fn read_xl_meta_no_data<R: AsyncRead + Unpin>(reader: &mut R, size: us
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod test {
|
mod test {
|
||||||
|
|
||||||
use std::fs;
|
|
||||||
use std::fs::File;
|
|
||||||
use std::os::unix::fs::MetadataExt;
|
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
@@ -123,6 +123,8 @@ pub struct MetacacheReader<R> {
|
|||||||
err: Option<Error>,
|
err: Option<Error>,
|
||||||
buf: Vec<u8>,
|
buf: Vec<u8>,
|
||||||
offset: usize,
|
offset: usize,
|
||||||
|
|
||||||
|
current: Option<MetaCacheEntry>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
||||||
@@ -133,6 +135,7 @@ impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
err: None,
|
err: None,
|
||||||
buf: Vec::new(),
|
buf: Vec::new(),
|
||||||
offset: 0,
|
offset: 0,
|
||||||
|
current: None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -199,7 +202,7 @@ impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
Marker::Str8 => Ok(u32::from(self.read_u8().await?)),
|
Marker::Str8 => Ok(u32::from(self.read_u8().await?)),
|
||||||
Marker::Str16 => Ok(u32::from(self.read_u16().await?)),
|
Marker::Str16 => Ok(u32::from(self.read_u16().await?)),
|
||||||
Marker::Str32 => Ok(self.read_u32().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<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
Ok(u32::from_be_bytes(buf.try_into().expect("Slice with incorrect length")))
|
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<Option<MetaCacheEntry>> {
|
pub async fn peek(&mut self) -> Result<Option<MetaCacheEntry>> {
|
||||||
self.check_init().await?;
|
self.check_init().await?;
|
||||||
|
|
||||||
@@ -279,12 +321,15 @@ impl<R: AsyncRead + Unpin> MetacacheReader<R> {
|
|||||||
|
|
||||||
self.reset();
|
self.reset();
|
||||||
|
|
||||||
Ok(Some(MetaCacheEntry {
|
let entry = Some(MetaCacheEntry {
|
||||||
name,
|
name,
|
||||||
metadata,
|
metadata,
|
||||||
cached: None,
|
cached: None,
|
||||||
reusable: false,
|
reusable: false,
|
||||||
}))
|
});
|
||||||
|
self.current = entry.clone();
|
||||||
|
|
||||||
|
Ok(entry)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn read_all(&mut self) -> Result<Vec<MetaCacheEntry>> {
|
pub async fn read_all(&mut self) -> Result<Vec<MetaCacheEntry>> {
|
||||||
|
|||||||
@@ -290,3 +290,132 @@ impl ECStore {
|
|||||||
Ok(ress)
|
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<Error>]| {
|
||||||
|
Box::pin(async move { println!("get entries: {:?}", entries) })
|
||||||
|
})),
|
||||||
|
finished: None,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user