From bfe8c5e9bf23dd7fc04b007051e0c1e932cabfdc Mon Sep 17 00:00:00 2001 From: weisd Date: Sun, 12 Jan 2025 22:36:57 +0800 Subject: [PATCH] add walk api, fix: load iam config --- ecstore/src/cache_value/metacache_set.rs | 4 +- ecstore/src/store_list_objects.rs | 755 ++++++++++++++++------- iam/src/manager.rs | 2 +- iam/src/store/object.rs | 112 ++-- rustfs/src/admin/handlers/user.rs | 4 - rustfs/src/auth.rs | 7 + 6 files changed, 606 insertions(+), 278 deletions(-) diff --git a/ecstore/src/cache_value/metacache_set.rs b/ecstore/src/cache_value/metacache_set.rs index 1f5535651..263b7fa08 100644 --- a/ecstore/src/cache_value/metacache_set.rs +++ b/ecstore/src/cache_value/metacache_set.rs @@ -54,7 +54,7 @@ impl Clone for ListPathRawOptions { } pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) -> Result<()> { - info!("list_path_raw"); + // println!("list_path_raw {},{}", &opts.bucket, &opts.path); if opts.disks.is_empty() { info!("list_path_raw 0 drives provided"); return Err(Error::from_string("list_path_raw: 0 drives provided")); @@ -283,11 +283,9 @@ pub async fn list_path_raw(mut rx: B_Receiver, opts: ListPathRawOptions) - } } - info!("read entry should heal: {}", current.name); if let Some(partial_fn) = opts.partial.as_ref() { partial_fn(MetaCacheEntries(top_entries), &errs).await; } - // break; } Ok(()) }); diff --git a/ecstore/src/store_list_objects.rs b/ecstore/src/store_list_objects.rs index 6bc169987..5d5bfa1e6 100644 --- a/ecstore/src/store_list_objects.rs +++ b/ecstore/src/store_list_objects.rs @@ -1,3 +1,5 @@ +use crate::bucket::metadata_sys::get_versioning_config; +use crate::bucket::versioning::VersioningApi; use crate::cache_value::metacache_set::{list_path_raw, ListPathRawOptions}; use crate::disk::error::{is_all_not_found, is_all_volume_not_found, is_err_eof, DiskError}; use crate::disk::{ @@ -9,10 +11,10 @@ use crate::file_meta::merge_file_meta_versions; use crate::peer::is_reserved_or_invalid_bucket; use crate::set_disk::SetDisks; use crate::store::check_list_objs_args; -use crate::store_api::{ListObjectVersionsInfo, ListObjectsInfo, ObjectInfo, ObjectOptions}; +use crate::store_api::{FileInfo, ListObjectVersionsInfo, ListObjectsInfo, ObjectInfo, ObjectOptions}; use crate::store_err::{is_err_bucket_not_found, to_object_err, StorageError}; use crate::utils::path::{self, base_dir_from_prefix, SLASH_SEPARATOR}; -use crate::StorageAPI; +use crate::{bucket, StorageAPI}; use crate::{store::ECStore, store_api::ListObjectsV2Info}; use futures::future::join_all; use rand::seq::SliceRandom; @@ -22,7 +24,7 @@ use std::io::ErrorKind; use std::sync::Arc; use tokio::sync::broadcast::{self, Receiver as B_Receiver}; use tokio::sync::mpsc::{self, Receiver, Sender}; -use tracing::error; +use tracing::{error, warn}; use uuid::Uuid; const MAX_OBJECT_LIST: i32 = 1000; @@ -245,7 +247,7 @@ impl ECStore { ..Default::default() }; - // warn!("list_objects_generic opts {:?}", &opts); + warn!("list_objects_generic opts {:?}", &opts); // use get if !opts.prefix.is_empty() && opts.limit == 1 && opts.marker.is_none() { @@ -669,6 +671,287 @@ impl ECStore { Ok(Vec::new()) } + + pub async fn walk( + self: Arc, + rx: B_Receiver, + bucket: &str, + prefix: &str, + result: Sender, + opts: WalkOptions, + ) -> Result<()> { + check_list_objs_args(bucket, prefix, &None)?; + + let mut futures = Vec::new(); + let mut inputs = Vec::new(); + + for eset in self.pools.iter() { + for set in eset.disk_set.iter() { + let (mut disks, infos, _) = set.get_online_disks_with_healing_and_info(true).await; + let rx = rx.resubscribe(); + let opts = opts.clone(); + + let (sender, list_out_rx) = mpsc::channel::(1); + inputs.push(list_out_rx); + futures.push(async move { + let mut ask_disks = get_list_quorum(&opts.ask_disks, set.set_drive_count as i32); + if ask_disks == -1 { + let new_disks = get_quorum_disks(&disks, &infos, (disks.len() + 1) / 2); + if !new_disks.is_empty() { + disks = new_disks; + } else { + ask_disks = get_list_quorum("strict", set.set_drive_count as i32); + } + } + + if set.set_drive_count == 4 || ask_disks > disks.len() as i32 { + ask_disks = disks.len() as i32; + } + + let fallback_disks = { + if ask_disks > 0 && disks.len() > ask_disks as usize { + let mut rand = thread_rng(); + disks.shuffle(&mut rand); + disks.split_off(ask_disks as usize) + } else { + Vec::new() + } + }; + + let listing_quorum = ((ask_disks + 1) / 2) as usize; + + let resolver = MetadataResolutionParams { + dir_quorum: listing_quorum, + obj_quorum: listing_quorum, + bucket: bucket.to_owned(), + ..Default::default() + }; + + let path = base_dir_from_prefix(prefix); + + let mut filter_prefix = { + prefix + .trim_start_matches(&path) + .trim_start_matches(SLASH_SEPARATOR) + .trim_end_matches(SLASH_SEPARATOR) + .to_owned() + }; + + if filter_prefix == path { + filter_prefix = "".to_owned(); + } + + // let (sender, rx1) = mpsc::channel(100); + + let tx1 = sender.clone(); + let tx2 = sender.clone(); + + list_path_raw( + rx.resubscribe(), + ListPathRawOptions { + disks: disks.iter().cloned().map(Some).collect(), + fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), + bucket: bucket.to_owned(), + path, + recursice: true, + filter_prefix: Some(filter_prefix), + forward_to: opts.marker.clone(), + min_disks: listing_quorum, + per_disk_limit: opts.limit as i32, + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + Box::pin({ + let value = tx1.clone(); + async move { + if entry.is_dir() { + return; + } + if let Err(err) = value.send(entry).await { + error!("list_path send fail {:?}", err); + } + } + }) + })), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + Box::pin({ + let value = tx2.clone(); + let resolver = resolver.clone(); + async move { + if let Ok(Some(entry)) = entries.resolve(resolver) { + if let Err(err) = value.send(entry).await { + error!("list_path send fail {:?}", err); + } + } + } + }) + })), + finished: None, + ..Default::default() + }, + ) + .await + }); + } + } + + let (merge_tx, mut merge_rx) = mpsc::channel::(100); + + let bucket = bucket.to_owned(); + + let vcf = match get_versioning_config(&bucket).await { + Ok((res, _)) => Some(res), + Err(_) => None, + }; + + tokio::spawn(async move { + let mut sent_err = false; + while let Some(entry) = merge_rx.recv().await { + if opts.latest_only { + let fi = match entry.to_fileinfo(&bucket) { + Ok(res) => res, + Err(err) => { + if !sent_err { + let item = ObjectInfoOrErr { + item: None, + err: Some(err), + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + + sent_err = true; + + return; + } + + continue; + } + }; + + if let Some(fiter) = opts.filter { + if fiter(&fi) { + let item = ObjectInfoOrErr { + item: Some(fi.to_object_info(&bucket, &fi.name, { + if let Some(v) = &vcf { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + } + } else { + let item = ObjectInfoOrErr { + item: Some(fi.to_object_info(&bucket, &fi.name, { + if let Some(v) = &vcf { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + } + continue; + } + + let fvs = match entry.file_info_versions(&bucket) { + Ok(res) => res, + Err(err) => { + let item = ObjectInfoOrErr { + item: None, + err: Some(err), + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + return; + } + }; + + if opts.versions_sort == WalkVersionsSortOrder::Ascending { + //TODO: SORT + } + + for fi in fvs.versions.iter() { + if let Some(fiter) = opts.filter { + if fiter(fi) { + let item = ObjectInfoOrErr { + item: Some(fi.to_object_info(&bucket, &fi.name, { + if let Some(v) = &vcf { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + } + } else { + let item = ObjectInfoOrErr { + item: Some(fi.to_object_info(&bucket, &fi.name, { + if let Some(v) = &vcf { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + } + } + } + }); + + tokio::spawn(async move { merge_entry_channels(rx, inputs, merge_tx, 1).await }); + + join_all(futures).await; + + Ok(()) + } +} + +type WalkFilter = fn(&FileInfo) -> bool; + +#[derive(Clone, Default)] +pub struct WalkOptions { + pub filter: Option, // return WalkFilter returns 'true/false' + pub marker: Option, // set to skip until this object + pub latest_only: bool, // returns only latest versions for all matching objects + pub ask_disks: String, // dictates how many disks are being listed + pub versions_sort: WalkVersionsSortOrder, // sort order for versions of the same object; default: Ascending order in ModTime + pub limit: usize, // maximum number of items, 0 means no limit +} + +#[derive(Clone, Default, PartialEq, Eq)] +pub enum WalkVersionsSortOrder { + #[default] + Ascending, + Descending, +} + +#[derive(Debug)] +pub struct ObjectInfoOrErr { + pub item: Option, + pub err: Option, } async fn gather_results( @@ -1092,260 +1375,294 @@ fn calc_common_counter(infos: &[DiskInfo], read_quorum: usize) -> u64 { // list_path_raw -// #[cfg(test)] -// mod test { -// use std::sync::Arc; +#[cfg(test)] +mod test { + use std::sync::Arc; -// 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::format::FormatV3; -// 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::endpoints::EndpointServerPools; -// use crate::error::Error; -// use crate::metacache::writer::MetacacheReader; -// use crate::set_disk::SetDisks; -// use crate::store::ECStore; -// use crate::store_list_objects::ListPathOptions; -// use futures::future::join_all; -// use lock::namespace_lock::NsLockMap; -// use tokio::sync::broadcast; -// use tokio::sync::mpsc; -// use tokio::sync::RwLock; -// use uuid::Uuid; + 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::format::FormatV3; + 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::endpoints::EndpointServerPools; + use crate::error::Error; + use crate::metacache::writer::MetacacheReader; + use crate::set_disk::SetDisks; + use crate::store::ECStore; + use crate::store_list_objects::ListPathOptions; + use crate::store_list_objects::WalkOptions; + use crate::store_list_objects::WalkVersionsSortOrder; + use futures::future::join_all; + use lock::namespace_lock::NsLockMap; + use tokio::sync::broadcast; + use tokio::sync::mpsc; + use tokio::sync::RwLock; + use uuid::Uuid; -// #[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; + // #[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 = 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 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 (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() -// }; + // 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); -// } -// }); + // 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); + // 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; -// } + // 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; -// } + // 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; + // #[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 = 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 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 = None; -// let disks = vec![Some(disk)]; -// let fallback_disks = Vec::new(); + // let (_, rx) = broadcast::channel(1); + // let bucket = "dada".to_owned(); + // let forward_to = None; + // 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(); -// } + // 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(); + // } -// #[tokio::test] -// async fn test_set_list_path() { -// 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; + // #[tokio::test] + // async fn test_set_list_path() { + // 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.set_disk_id(Some(Uuid::new_v4())).await; + // let disk = new_disk(&ep, &DiskOption::default()).await.expect("init disk fail"); + // let _ = disk.set_disk_id(Some(Uuid::new_v4())).await; -// let set = SetDisks { -// lockers: Vec::new(), -// locker_owner: String::new(), -// ns_mutex: Arc::new(RwLock::new(NsLockMap::new(false))), -// disks: RwLock::new(vec![Some(disk)]), -// set_endpoints: Vec::new(), -// set_drive_count: 1, -// default_parity_count: 0, -// set_index: 0, -// pool_index: 0, -// format: FormatV3::new(1, 1), -// }; + // let set = SetDisks { + // lockers: Vec::new(), + // locker_owner: String::new(), + // ns_mutex: Arc::new(RwLock::new(NsLockMap::new(false))), + // disks: RwLock::new(vec![Some(disk)]), + // set_endpoints: Vec::new(), + // set_drive_count: 1, + // default_parity_count: 0, + // set_index: 0, + // pool_index: 0, + // format: FormatV3::new(1, 1), + // }; -// let (_tx, rx) = broadcast::channel(1); + // let (_tx, rx) = broadcast::channel(1); -// let bucket = "dada".to_owned(); + // let bucket = "dada".to_owned(); -// let opts = ListPathOptions { -// bucket, -// recursive: true, -// ..Default::default() -// }; + // let opts = ListPathOptions { + // bucket, + // recursive: true, + // ..Default::default() + // }; -// let (sender, mut recv) = mpsc::channel(10); + // let (sender, mut recv) = mpsc::channel(10); -// set.list_path(rx, opts, sender).await.unwrap(); + // set.list_path(rx, opts, sender).await.unwrap(); -// while let Some(entry) = recv.recv().await { -// println!("get entry {:?}", entry.name) -// } -// } + // while let Some(entry) = recv.recv().await { + // println!("get entry {:?}", entry.name) + // } + // } -// #[tokio::test] -// async fn test_list_merged() { -// let server_address = "localhost:9000"; + // #[tokio::test] + //walk() { + // let server_address = "localhost:9000"; -// let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( -// server_address, -// vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], -// ) -// .unwrap(); + // let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( + // server_address, + // vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], + // ) + // .unwrap(); -// let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) -// .await -// .unwrap(); + // let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) + // .await + // .unwrap(); -// let (_tx, rx) = broadcast::channel(1); + // let (_tx, rx) = broadcast::channel(1); -// let bucket = "dada".to_owned(); -// let opts = ListPathOptions { -// bucket, -// recursive: true, -// ..Default::default() -// }; + // let bucket = "dada".to_owned(); + // let opts = ListPathOptions { + // bucket, + // recursive: true, + // ..Default::default() + // }; -// let (sender, mut recv) = mpsc::channel(10); + // let (sender, mut recv) = mpsc::channel(10); -// store.list_merged(rx, opts, sender).await.unwrap(); + // store.list_merged(rx, opts, sender).await.unwrap(); -// while let Some(entry) = recv.recv().await { -// println!("get entry {:?}", entry.name) -// } -// } + // while let Some(entry) = recv.recv().await { + // println!("get entry {:?}", entry.name) + // } + // } -// #[tokio::test] -// async fn test_list_path() { -// let server_address = "localhost:9000"; + // #[tokio::test] + // async fn test_list_path() { + // let server_address = "localhost:9000"; -// let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( -// server_address, -// vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], -// ) -// .unwrap(); + // let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( + // server_address, + // vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], + // ) + // .unwrap(); -// let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) -// .await -// .unwrap(); + // let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) + // .await + // .unwrap(); -// let bucket = "dada".to_owned(); -// let opts = ListPathOptions { -// bucket, -// recursive: true, -// limit: 100, + // let bucket = "dada".to_owned(); + // let opts = ListPathOptions { + // bucket, + // recursive: true, + // limit: 100, -// ..Default::default() -// }; + // ..Default::default() + // }; -// let ret = store.list_path(&opts).await.unwrap(); -// println!("ret {:?}", ret); -// } + // let ret = store.list_path(&opts).await.unwrap(); + // println!("ret {:?}", ret); + // } -// // #[tokio::test] -// // async fn test_list_objects_v2() { -// // let server_address = "localhost:9000"; + // #[tokio::test] + // async fn test_list_objects_v2() { + // let server_address = "localhost:9000"; -// // let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( -// // server_address, -// // vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], -// // ) -// // .unwrap(); + // let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( + // server_address, + // vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], + // ) + // .unwrap(); -// // let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) -// // .await -// // .unwrap(); + // let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) + // .await + // .unwrap(); -// // let ret = store.list_objects_v2("data", "", "", "", 100, false, "").await.unwrap(); -// // println!("ret {:?}", ret); -// // } -// } + // let ret = store.list_objects_v2("data", "", "", "", 100, false, "").await.unwrap(); + // println!("ret {:?}", ret); + // } + + #[tokio::test] + async fn test_walk() { + let server_address = "localhost:9000"; + + let (endpoint_pools, _setup_type) = EndpointServerPools::from_volumes( + server_address, + vec!["/Users/weisd/project/weisd/s3-rustfs/target/volume/test".to_string()], + ) + .unwrap(); + + let store = ECStore::new(server_address.to_string(), endpoint_pools.clone()) + .await + .unwrap(); + + ECStore::init(store.clone()).await.unwrap(); + + let (_tx, rx) = broadcast::channel(1); + + let bucket = ".rustfs.sys"; + let prefix = "config/iam/sts/"; + + let (sender, mut recv) = mpsc::channel(10); + + let opts = WalkOptions::default(); + + store.walk(rx, bucket, prefix, sender, opts).await.unwrap(); + + while let Some(entry) = recv.recv().await { + println!("get entry {:?}", entry) + } + } +} diff --git a/iam/src/manager.rs b/iam/src/manager.rs index 0be502209..386893e13 100644 --- a/iam/src/manager.rs +++ b/iam/src/manager.rs @@ -261,7 +261,7 @@ where self.api.save_iam_config(&user_entiry, path).await?; Cache::add_or_update( - &self.cache.users, + &self.cache.sts_accounts, &user_entiry.credentials.access_key, &user_entiry, OffsetDateTime::now_utc(), diff --git a/iam/src/store/object.rs b/iam/src/store/object.rs index 9290395ff..fd4fadc24 100644 --- a/iam/src/store/object.rs +++ b/iam/src/store/object.rs @@ -4,12 +4,14 @@ use ecstore::{ config::error::is_err_config_not_found, store::ECStore, store_api::{ObjectIO, ObjectInfo, ObjectOptions, PutObjReader}, - utils::path::dir, - StorageAPI, + store_list_objects::{ObjectInfoOrErr, WalkOptions}, + utils::path::{dir, SLASH_SEPARATOR}, }; use futures::future::try_join_all; use log::{debug, warn}; use serde::{de::DeserializeOwned, Serialize}; +use tokio::sync::broadcast; +use tokio::sync::mpsc::{self}; use super::Store; use crate::{ @@ -31,45 +33,54 @@ impl ObjectStore { Self { object_api } } - async fn list_iam_config_items(&self, prefix: &str, items: &[&str]) -> crate::Result> { - debug!("list iam config items, prefix: {prefix}"); + async fn list_iam_config_items(&self, prefix: &str) -> crate::Result> { + debug!("list iam config items, prefix: {}", &prefix); // todo, 实现walk,使用walk - let mut futures = Vec::with_capacity(items.len()); - for item in items { - let prefix = format!("{}{}", prefix, item); - futures.push(async move { - // let items = self - // .object_api - // .clone() - // .list_path(&ListPathOptions { - // bucket: Self::BUCKET_NAME.into(), - // prefix: prefix.clone(), - // ..Default::default() - // }) - // .await; + let (ctx_tx, ctx_rx) = broadcast::channel(1); - let items = self - .object_api - .clone() - .list_objects_v2(Self::BUCKET_NAME, &prefix.clone(), None, None, 0, false, None) - .await; + // let prefix = format!("{}{}", prefix, item); - match items { - Ok(items) => Result::<_, crate::Error>::Ok(items.prefixes), - Err(e) => { - if is_err_config_not_found(&e) { - Result::<_, crate::Error>::Ok(vec![]) - } else { - Err(Error::StringError(format!("list {prefix} failed, err: {e:?}"))) - } - } - } - }); + let ctx_rx = ctx_rx.resubscribe(); + + let (tx, mut rx) = mpsc::channel::(100); + let store = self.object_api.clone(); + + let path = prefix.to_owned(); + tokio::spawn(async move { store.walk(ctx_rx, Self::BUCKET_NAME, &path, tx, WalkOptions::default()).await }); + + let mut ret = Vec::new(); + + while let Some(v) = rx.recv().await { + if let Some(err) = v.err { + warn!("list_iam_config_items {:?}", err); + let _ = ctx_tx.send(true); + + return Err(Error::EcstoreError(err)); + } + + if let Some(item) = v.item { + let name = item.name.trim_start_matches(prefix).trim_end_matches(SLASH_SEPARATOR); + ret.push(name.to_owned()); + } } + + Ok(ret) + + // match items { + // Ok(items) => Result::<_, crate::Error>::Ok(items.prefixes), + // Err(e) => { + // if is_err_config_not_found(&e) { + // Result::<_, crate::Error>::Ok(vec![]) + // } else { + // Err(Error::StringError(format!("list {prefix} failed, err: {e:?}"))) + // } + // } + // } + // TODO: FIXME: - Ok(try_join_all(futures).await?.into_iter().flat_map(|x| x.into_iter()).collect()) + // Ok(try_join_all(futures).await?.into_iter().flat_map(|x| x.into_iter()).collect()) } async fn load_policy(&self, name: &str) -> crate::Result { @@ -152,7 +163,7 @@ impl Store for ObjectStore { } async fn load_policy_docs(&self) -> crate::Result> { - let paths = self.list_iam_config_items("config/iam/", &["policies/"]).await?; + let paths = self.list_iam_config_items("config/iam/policies/").await?; let mut result = Self::get_default_policyes(); for path in paths { @@ -174,7 +185,9 @@ impl Store for ObjectStore { } async fn load_users(&self, user_type: UserType) -> crate::Result> { - let paths = self.list_iam_config_items("config/iam/", &[user_type.prefix()]).await?; + let paths = self + .list_iam_config_items(format!("{}{}", "config/iam/", user_type.prefix()).as_str()) + .await?; let mut result = HashMap::new(); for path in paths { @@ -202,21 +215,18 @@ impl Store for ObjectStore { /// load all and make a new cache. async fn load_all(&self, cache: &Cache) -> crate::Result<()> { - let items = self - .list_iam_config_items( - "config/iam/", - &[ - "policydb/", - "policies/", - "groups/", - "policydb/users/", - "policydb/groups/", - "service-accounts/", - "policydb/sts-users/", - "sts/", - ], - ) - .await?; + let _items = &[ + "policydb/", + "policies/", + "groups/", + "policydb/users/", + "policydb/groups/", + "service-accounts/", + "policydb/sts-users/", + "sts/", + ]; + + let items = self.list_iam_config_items("config/iam/").await?; debug!("all iam items: {items:?}"); let (policy_docs, users, user_policies, sts_policies, sts_accounts) = ( diff --git a/rustfs/src/admin/handlers/user.rs b/rustfs/src/admin/handlers/user.rs index fe9dd8356..63c24895b 100644 --- a/rustfs/src/admin/handlers/user.rs +++ b/rustfs/src/admin/handlers/user.rs @@ -66,10 +66,6 @@ impl Operation for AddUser { } } -fn check_claims_from_token(_token: &str) -> bool { - true -} - pub struct SetUserStatus {} #[async_trait::async_trait] impl Operation for SetUserStatus { diff --git a/rustfs/src/auth.rs b/rustfs/src/auth.rs index f76fdfcbb..ee6678c7d 100644 --- a/rustfs/src/auth.rs +++ b/rustfs/src/auth.rs @@ -1,4 +1,5 @@ use iam::cache::CacheInner; +use log::warn; use s3s::auth::S3Auth; use s3s::auth::SecretKey; use s3s::auth::SimpleAuth; @@ -27,9 +28,15 @@ impl S3Auth for IAMAuth { return Ok(key); } + warn!("Failed to get secret key from simple auth"); + if let Ok(iam_store) = iam::get() { let c = CacheInner::from(&iam_store.cache); + warn!("Failed to get secret key from simple auth, try cache {}", access_key); + warn!("users {:?}", c.users.values()); + warn!("sts_accounts {:?}", c.sts_accounts.values()); if let Some(id) = c.get_user(access_key) { + warn!("get cred {:?}", id.credentials); return Ok(SecretKey::from(id.credentials.secret_key.clone())); } }