From 583a23bdf29597e3493e4f078f392508c469278d Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 23 Jun 2026 12:31:17 +0800 Subject: [PATCH] fix(ecstore): replace panic-driven pool and set stubs (#3753) * fix(ecstore): replace panic-driven pool and set stubs * test(runtime): tolerate restricted local bind checks * fix(ecstore): remove remaining trait stub placeholders * fix(ecstore): tighten trait stub follow-up semantics * chore: ignore local worktrees * chore: update layer dependency baseline for resolve_* context entries Add 7 accepted infra->app dependency entries introduced by recent refactoring PRs (#3770, #3771, #3772) that route global state lookups through app::context::resolve_* functions. Co-Authored-By: heihutu --------- Co-authored-by: heihutu --- .gitignore | 3 +- crates/ecstore/src/client/api_remove.rs | 19 +- crates/ecstore/src/rpc/mod.rs | 1 + crates/ecstore/src/rpc/peer_s3_client.rs | 6 +- crates/ecstore/src/rpc/remote_disk.rs | 24 +- crates/ecstore/src/rpc/remote_locker.rs | 20 +- crates/ecstore/src/set_disk.rs | 580 ++++++++- crates/ecstore/src/sets.rs | 316 ++++- crates/ecstore/src/store/heal.rs | 25 +- crates/ecstore/src/store/multipart.rs | 5 +- crates/ecstore/src/store_list_objects.rs | 1302 ++++++++++++++++++++ crates/iam/src/oidc.rs | 20 +- crates/rio/src/http_reader.rs | 40 +- crates/targets/src/target/redis.rs | 30 +- crates/utils/src/net.rs | 58 +- rustfs/src/startup_server.rs | 37 +- rustfs/src/storage/rpc/node_service.rs | 24 +- rustfs/tests/embedded_deferred_iam_test.rs | 6 +- rustfs/tests/embedded_test.rs | 6 +- scripts/layer-dependency-baseline.txt | 1 + 20 files changed, 2361 insertions(+), 162 deletions(-) diff --git a/.gitignore b/.gitignore index 622505974..21fb4f7cc 100644 --- a/.gitignore +++ b/.gitignore @@ -65,4 +65,5 @@ benchmarks.logs tmp/ crates/*/docs fuzz/target -outputs \ No newline at end of file +outputs +worktrees/* \ No newline at end of file diff --git a/crates/ecstore/src/client/api_remove.rs b/crates/ecstore/src/client/api_remove.rs index 8164eb962..f192087da 100644 --- a/crates/ecstore/src/client/api_remove.rs +++ b/crates/ecstore/src/client/api_remove.rs @@ -765,9 +765,16 @@ mod tests { net::TcpListener, }; - async fn capture_delete_objects_sha256_header() -> (String, tokio::task::JoinHandle) { - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let endpoint = listener.local_addr().unwrap().to_string(); + async fn capture_delete_objects_sha256_header() -> Option<(String, tokio::task::JoinHandle)> { + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, + Err(err) => panic!("test listener should bind: {err}"), + }; + let endpoint = listener + .local_addr() + .expect("listener local address should be available") + .to_string(); let task = tokio::spawn(async move { let (mut stream, _) = listener.accept().await.unwrap(); let mut request = Vec::new(); @@ -801,7 +808,7 @@ mod tests { sha256_header }); - (endpoint, task) + Some((endpoint, task)) } #[tokio::test] @@ -813,7 +820,9 @@ mod tests { }]; let body = generate_remove_multi_objects_request(&objects); let expected = rustfs_utils::hex(HashAlgorithm::SHA256.hash_encode(&body)); - let (endpoint, header_task) = capture_delete_objects_sha256_header().await; + let Some((endpoint, header_task)) = capture_delete_objects_sha256_header().await else { + return; + }; let client = TransitionClient::new( &endpoint, Options { diff --git a/crates/ecstore/src/rpc/mod.rs b/crates/ecstore/src/rpc/mod.rs index a862b6777..253eaeaa8 100644 --- a/crates/ecstore/src/rpc/mod.rs +++ b/crates/ecstore/src/rpc/mod.rs @@ -31,6 +31,7 @@ pub use internode_data_transport::build_internode_data_transport_from_env; pub use peer_rest_client::{ PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, }; +pub(crate) use peer_s3_client::heal_bucket_local_on_disks; pub use peer_s3_client::{LocalPeerS3Client, PeerS3Client, S3PeerSys}; pub use remote_disk::RemoteDisk; pub use remote_locker::RemoteClient; diff --git a/crates/ecstore/src/rpc/peer_s3_client.rs b/crates/ecstore/src/rpc/peer_s3_client.rs index 8ad31466b..e5bae079c 100644 --- a/crates/ecstore/src/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/rpc/peer_s3_client.rs @@ -954,7 +954,11 @@ pub async fn heal_bucket_local(bucket: &str, opts: &HealOpts) -> Result>) -> Result { +pub(crate) async fn heal_bucket_local_on_disks( + bucket: &str, + opts: &HealOpts, + disks: Vec>, +) -> Result { let before_state = Arc::new(RwLock::new(vec![String::new(); disks.len()])); let after_state = Arc::new(RwLock::new(vec![String::new(); disks.len()])); diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 285bdeab0..6923a9117 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -2662,8 +2662,12 @@ mod tests { #[tokio::test] async fn test_remote_disk_is_online_detects_active_listener() { - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); let url = url::Url::parse(&format!("http://{}:{}/data/rustfs0", addr.ip(), addr.port())).unwrap(); let endpoint = Endpoint { @@ -2691,8 +2695,12 @@ mod tests { async fn test_remote_disk_is_online_detects_missing_listener() { init_tracing(Level::ERROR); - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); let ip = addr.ip(); let port = addr.port(); @@ -2747,8 +2755,12 @@ mod tests { init_tracing(Level::ERROR); let _ = rustfs_credentials::GLOBAL_RUSTFS_RPC_SECRET.set("test-rpc-secret".to_string()); - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); let accept_task = tokio::spawn(async move { while let Ok((stream, _)) = listener.accept().await { drop(stream); diff --git a/crates/ecstore/src/rpc/remote_locker.rs b/crates/ecstore/src/rpc/remote_locker.rs index 3138920d4..8f616da34 100644 --- a/crates/ecstore/src/rpc/remote_locker.rs +++ b/crates/ecstore/src/rpc/remote_locker.rs @@ -558,16 +558,20 @@ mod tests { use tokio::task::JoinHandle; use tonic::transport::Endpoint as TonicEndpoint; - async fn spawn_hanging_listener() -> (String, JoinHandle<()>) { - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = format!("http://{}", listener.local_addr().unwrap()); + async fn spawn_hanging_listener() -> Option<(String, JoinHandle<()>)> { + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = format!("http://{}", listener.local_addr().expect("listener local address should be available")); let task = tokio::spawn(async move { if let Ok((stream, _)) = listener.accept().await { let _stream = stream; tokio::time::sleep(Duration::from_secs(2)).await; } }); - (addr, task) + Some((addr, task)) } async fn cache_lazy_channel(addr: &str) { @@ -588,7 +592,9 @@ mod tests { #[tokio::test] async fn test_remote_client_acquire_lock_respects_request_timeout_and_evicts_connection() { ensure_test_rpc_secret(); - let (addr, accept_task) = spawn_hanging_listener().await; + let Some((addr, accept_task)) = spawn_hanging_listener().await else { + return; + }; cache_lazy_channel(&addr).await; assert!(GLOBAL_CONN_MAP.read().await.contains_key(&addr)); @@ -622,7 +628,9 @@ mod tests { #[tokio::test] async fn test_remote_client_acquire_locks_batch_respects_request_timeout_and_evicts_connection() { ensure_test_rpc_secret(); - let (addr, accept_task) = spawn_hanging_listener().await; + let Some((addr, accept_task)) = spawn_hanging_listener().await else { + return; + }; cache_lazy_channel(&addr).await; assert!(GLOBAL_CONN_MAP.read().await.contains_key(&addr)); diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index cd02df998..172991888 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -18,12 +18,15 @@ use crate::batch_processor::{AsyncBatchProcessor, get_global_processors}; use crate::bitrot::{create_bitrot_reader, create_bitrot_writer}; use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE; +use crate::bucket::metadata_sys; use crate::bucket::object_lock::objectlock_sys::check_retention_for_modification; use crate::bucket::replication::check_replicate_delete; use crate::bucket::versioning::VersioningApi; use crate::bucket::versioning_sys::BucketVersioningSys; use crate::client::{object_api_utils::get_raw_etag, transition_api::ReaderImpl}; -use crate::disk::error_reduce::{OBJECT_OP_IGNORED_ERRS, reduce_read_quorum_errs, reduce_write_quorum_errs}; +use crate::disk::error_reduce::{ + BUCKET_OP_IGNORED_ERRS, OBJECT_OP_IGNORED_ERRS, count_errs, reduce_read_quorum_errs, reduce_write_quorum_errs, +}; use crate::disk::{ self, CHECK_PART_DISK_NOT_FOUND, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN, conv_part_err_to_int, has_part_err, @@ -34,6 +37,8 @@ use crate::error::{Error, Result, is_err_version_not_found}; use crate::error::{GenericError, ObjectApiError, is_err_object_not_found}; use crate::global::{GLOBAL_LocalNodeName, GLOBAL_TierConfigMgr}; use crate::object_api::ObjectOptions; +use crate::rpc::heal_bucket_local_on_disks; +use crate::store_utils::is_reserved_or_invalid_bucket; use crate::{ bucket::lifecycle::bucket_lifecycle_ops::{ LifecycleOps, gen_transition_objname, get_transitioned_object_reader, put_restore_opts, @@ -50,7 +55,7 @@ use crate::{ event_notification::{EventArgs, send_event}, global::{GLOBAL_LOCAL_DISK_MAP, GLOBAL_LOCAL_DISK_SET_DRIVES, get_global_deployment_id, is_dist_erasure}, object_api::{GetObjectReader, ObjectInfo, PutObjReader}, - store_init::load_format_erasure, + store_init::{get_format_erasure_in_quorum, load_format_erasure, load_format_erasure_all, save_format_file}, }; use bytes::Bytes; use bytesize::ByteSize; @@ -75,7 +80,7 @@ use rustfs_lock::LockClient; use rustfs_lock::fast_lock::types::LockResult; use rustfs_lock::local_lock::LocalLock; use rustfs_lock::{FastLockGuard, LockManager, NamespaceLock, NamespaceLockGuard, NamespaceLockWrapper, ObjectKey}; -use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem}; +use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos}; use rustfs_object_capacity::capacity_scope::{ CapacityScope, CapacityScopeDisk, record_capacity_scope, record_global_dirty_scope, }; @@ -1772,23 +1777,205 @@ impl BucketOperations for SetDisks { type Error = Error; #[tracing::instrument(skip(self))] - async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> { - unimplemented!() + async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> { + let disks = self.disk_inventory().await; + let write_quorum = (disks.len() / 2) + 1; + let force_create = opts.force_create; + + let mut futures = Vec::with_capacity(disks.len()); + for disk in disks { + let bucket = bucket.to_string(); + futures.push(async move { + match disk { + Some(disk) => match disk.make_volume(&bucket).await { + Ok(()) => Ok(()), + Err(err) if force_create && matches!(err, DiskError::VolumeExists) => Ok(()), + Err(err) => Err(err), + }, + None => Err(DiskError::DiskNotFound), + } + }); + } + + let results = join_all(futures).await; + let errs = results + .into_iter() + .map(|result| result.err()) + .collect::>>(); + + if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) { + return Err(err.into()); + } + + Ok(()) } #[tracing::instrument(skip(self))] - async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { - unimplemented!() + async fn get_bucket_info(&self, bucket: &str, _opts: &BucketOptions) -> Result { + let disks = self.disk_inventory().await; + let write_quorum = (disks.len() / 2) + 1; + + let mut futures = Vec::with_capacity(disks.len()); + for disk in disks { + let bucket = bucket.to_string(); + futures.push(async move { + match disk { + Some(disk) => disk.stat_volume(&bucket).await, + None => Err(DiskError::DiskNotFound), + } + }); + } + + let results = join_all(futures).await; + let mut infos = Vec::with_capacity(results.len()); + let mut errs = Vec::with_capacity(results.len()); + for result in results { + match result { + Ok(info) => { + infos.push(Some(info)); + errs.push(None); + } + Err(err) => { + infos.push(None); + errs.push(Some(err)); + } + } + } + + if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) { + return Err(err.into()); + } + + let mut versioning = false; + let mut object_locking = false; + if let Ok(sys) = metadata_sys::get(bucket).await { + versioning = sys.versioning(); + object_locking = sys.object_locking(); + } + + infos + .into_iter() + .flatten() + .next() + .map(|info| BucketInfo { + name: info.name, + created: info.created, + versioning, + object_locking, + ..Default::default() + }) + .ok_or(Error::VolumeNotFound) } #[tracing::instrument(skip(self))] async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { - unimplemented!() + let disks = self.disk_inventory().await; + let write_quorum = (disks.len() / 2) + 1; + + let mut futures = Vec::with_capacity(disks.len()); + for disk in disks { + futures.push(async move { + match disk { + Some(disk) => disk.list_volumes().await, + None => Err(DiskError::DiskNotFound), + } + }); + } + + let results = join_all(futures).await; + let mut infos = Vec::with_capacity(results.len()); + let mut errs = Vec::with_capacity(results.len()); + for result in results { + match result { + Ok(volumes) => { + infos.push(Some(volumes)); + errs.push(None); + } + Err(err) => { + infos.push(None); + errs.push(Some(err)); + } + } + } + + if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) { + return Err(err.into()); + } + + let mut counts: HashMap = HashMap::new(); + for volumes in infos.into_iter().flatten() { + for volume in volumes { + if is_reserved_or_invalid_bucket(&volume.name, false) { + continue; + } + + let entry = counts.entry(volume.name.clone()).or_insert(( + 0, + BucketInfo { + name: volume.name.clone(), + created: volume.created, + ..Default::default() + }, + )); + entry.0 += 1; + } + } + + let mut buckets = counts + .into_values() + .filter_map(|(count, bucket)| (count >= write_quorum).then_some(bucket)) + .collect::>(); + buckets.sort_by(|left, right| left.name.cmp(&right.name)); + Ok(buckets) } #[tracing::instrument(skip(self))] - async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { - unimplemented!() + async fn delete_bucket(&self, bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { + let disks = self.disk_inventory().await; + let write_quorum = (disks.len() / 2) + 1; + + let mut futures = Vec::with_capacity(disks.len()); + for disk in disks.iter().cloned() { + let bucket = bucket.to_string(); + futures.push(async move { + match disk { + Some(disk) => disk.delete_volume(&bucket).await, + None => Err(DiskError::DiskNotFound), + } + }); + } + + let results = join_all(futures).await; + let mut errs = Vec::with_capacity(results.len()); + let mut recreate = false; + for result in results { + match result { + Ok(()) => errs.push(None), + Err(err) => { + if matches!(err, DiskError::VolumeNotEmpty) { + recreate = true; + } + errs.push(Some(err)); + } + } + } + + if recreate { + for (index, err) in errs.iter().enumerate() { + if err.is_none() + && let Some(Some(disk)) = disks.get(index) + { + let _ = disk.make_volume(bucket).await; + } + } + return Err(Error::VolumeNotEmpty); + } + + if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) { + return Err(err.into()); + } + + Ok(()) } } @@ -3027,40 +3214,51 @@ impl rustfs_storage_api::ListOperations for SetDisks { #[tracing::instrument(skip(self))] async fn list_objects_v2( self: Arc, - _bucket: &str, - _prefix: &str, - _continuation_token: Option, - _delimiter: Option, - _max_keys: i32, - _fetch_owner: bool, - _start_after: Option, - _incl_deleted: bool, + bucket: &str, + prefix: &str, + continuation_token: Option, + delimiter: Option, + max_keys: i32, + fetch_owner: bool, + start_after: Option, + incl_deleted: bool, ) -> Result { - unimplemented!() + self.inner_list_objects_v2( + bucket, + prefix, + continuation_token, + delimiter, + max_keys, + fetch_owner, + start_after, + incl_deleted, + ) + .await } #[tracing::instrument(skip(self))] async fn list_object_versions( self: Arc, - _bucket: &str, - _prefix: &str, - _marker: Option, - _version_marker: Option, - _delimiter: Option, - _max_keys: i32, + bucket: &str, + prefix: &str, + marker: Option, + version_marker: Option, + delimiter: Option, + max_keys: i32, ) -> Result { - unimplemented!() + self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys) + .await } async fn walk( self: Arc, - _rx: CancellationToken, - _bucket: &str, - _prefix: &str, - _result: Sender, - _opts: WalkOptions, + rx: CancellationToken, + bucket: &str, + prefix: &str, + result: Sender, + opts: WalkOptions, ) -> Result<()> { - unimplemented!() + self.walk_internal(rx, bucket, prefix, result, opts).await } } @@ -3092,7 +3290,7 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { _src_opts: &ObjectOptions, _dst_opts: &ObjectOptions, ) -> Result<()> { - unimplemented!() + Err(StorageError::NotImplemented) } #[tracing::instrument(level = "debug", skip(self, data, opts))] @@ -4155,13 +4353,67 @@ impl rustfs_storage_api::HealOperations for SetDisks { type HealOptions = HealOpts; #[tracing::instrument(skip(self))] - async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option)> { - unimplemented!() + async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { + let disks = self.disks.read().await.clone(); + let (formats, errs) = load_format_erasure_all(&disks, true).await; + let ref_format = match get_format_erasure_in_quorum(&formats) { + Ok(format) => format, + Err(err) => { + let can_use_cached_layout = count_errs(&errs, &DiskError::UnformattedDisk) > 0 + && formats.iter().flatten().all(|format| self.format.check_other(format).is_ok()) + && errs + .iter() + .all(|err| err.is_none() || matches!(err, Some(DiskError::UnformattedDisk))); + if can_use_cached_layout { + self.format.clone() + } else { + return Ok((HealResultItem::default(), Some(err))); + } + } + }; + + let endpoints = crate::endpoints::Endpoints::from(self.set_endpoints.clone()); + let before_drives = crate::layout::set_heal::formats_to_drives_info(&endpoints, &formats, &errs); + let mut result = HealResultItem { + heal_item_type: HealItemType::Metadata.to_string(), + detail: "disk-format".to_string(), + disk_count: self.set_drive_count, + set_count: 1, + before: Infos { + drives: before_drives.clone(), + }, + after: Infos { drives: before_drives }, + ..Default::default() + }; + + if count_errs(&errs, &DiskError::UnformattedDisk) == 0 { + info!("set disk formats success, NoHealRequired, errs: {:?}", errs); + return Ok((result, Some(StorageError::NoHealRequired))); + } + + if !dry_run { + for (disk_idx, err) in errs.iter().enumerate() { + if !matches!(err, Some(DiskError::UnformattedDisk)) { + continue; + } + + let mut new_format = ref_format.clone(); + new_format.erasure.this = ref_format.erasure.sets[self.set_index][disk_idx]; + if save_format_file(&disks[disk_idx], &Some(new_format.clone())).await.is_ok() { + result.after.drives[disk_idx].uuid = new_format.erasure.this.to_string(); + result.after.drives[disk_idx].state = DriveState::Ok.to_string(); + } + } + } + + Ok((result, None)) } #[tracing::instrument(skip(self))] - async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result { - unimplemented!() + async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result { + let mut result = heal_bucket_local_on_disks(bucket, opts, self.disk_inventory().await).await?; + result.set_count = 1; + Ok(result) } #[tracing::instrument(skip(self))] @@ -4237,13 +4489,23 @@ impl rustfs_storage_api::HealOperations for SetDisks { } #[tracing::instrument(skip(self))] - async fn get_pool_and_set(&self, _id: &str) -> Result<(Option, Option, Option)> { - unimplemented!() + async fn get_pool_and_set(&self, id: &str) -> Result<(Option, Option, Option)> { + for (set_idx, set) in self.format.erasure.sets.iter().enumerate() { + for (disk_idx, disk_id) in set.iter().enumerate() { + if disk_id.to_string() == id { + return Ok((Some(self.pool_index), Some(set_idx), Some(disk_idx))); + } + } + } + + Err(Error::DiskNotFound) } #[tracing::instrument(skip(self))] async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> { - unimplemented!() + // Multipart orphan reconciliation is intentionally retained above the set layer + // until there is a concrete caller and a stable lower-level contract to implement. + Err(StorageError::NotImplemented) } } @@ -5040,6 +5302,8 @@ mod tests { use rustfs_filemeta::ReplicationState; use rustfs_lock::client::local::LocalClient; use rustfs_lock::{LockError, LockInfo, LockResponse, LockStats}; + use rustfs_storage_api::HealOperations as _; + use rustfs_storage_api::ListOperations as _; use rustfs_storage_api::TransitionedObject; use rustfs_storage_api::{CompletePart, NamespaceLocking as _, ObjectOperations as _}; use serial_test::serial; @@ -7256,4 +7520,240 @@ mod tests { ); } } + + async fn make_local_bucket_test_set_disks() -> Arc { + let format = FormatV3::new(1, 2); + let mut endpoints = Vec::new(); + let mut disks = Vec::new(); + + for disk_idx in 0..2 { + let dir = tempfile::tempdir().expect("tempdir should be created"); + let mut endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(disk_idx); + + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("disk should be created"); + + let mut disk_format = format.clone(); + disk_format.erasure.this = format.erasure.sets[0][disk_idx]; + save_format_file(&Some(disk.clone()), &Some(disk_format)) + .await + .expect("format should be saved"); + + std::mem::forget(dir); + endpoints.push(endpoint); + disks.push(Some(disk)); + } + + SetDisks::new( + "test-owner".to_string(), + Arc::new(RwLock::new(disks)), + 2, + 1, + 0, + 0, + endpoints, + format, + Vec::new(), + ) + .await + } + + async fn make_local_bucket_test_set_disks_with_missing_format() -> Arc { + let format = FormatV3::new(1, 2); + let mut endpoints = Vec::new(); + let mut disks = Vec::new(); + + for disk_idx in 0..2 { + let dir = tempfile::tempdir().expect("tempdir should be created"); + let mut endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(disk_idx); + + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("disk should be created"); + + if disk_idx == 0 { + let mut disk_format = format.clone(); + disk_format.erasure.this = format.erasure.sets[0][disk_idx]; + save_format_file(&Some(disk.clone()), &Some(disk_format)) + .await + .expect("format should be saved"); + } + + std::mem::forget(dir); + endpoints.push(endpoint); + disks.push(Some(disk)); + } + + SetDisks::new( + "test-owner".to_string(), + Arc::new(RwLock::new(disks)), + 2, + 1, + 0, + 0, + endpoints, + format, + Vec::new(), + ) + .await + } + + #[tokio::test] + async fn bucket_operations_round_trip_without_panicking() { + let set_disks = make_local_bucket_test_set_disks().await; + let bucket = "bucket-roundtrip"; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + + let info = set_disks + .get_bucket_info(bucket, &BucketOptions::default()) + .await + .expect("bucket info should be available"); + assert_eq!(info.name, bucket); + + let buckets = set_disks + .list_bucket(&BucketOptions::default()) + .await + .expect("bucket listing should succeed"); + assert!(buckets.iter().any(|entry| entry.name == bucket)); + + set_disks + .delete_bucket(bucket, &DeleteBucketOptions::default()) + .await + .expect("bucket should be deleted"); + } + + #[tokio::test] + async fn set_level_listing_trait_methods_use_existing_listing_implementation() { + let set_disks = make_local_bucket_test_set_disks().await; + let bucket = "bucket-listing"; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + + let mut reader = PutObjReader::from_vec(b"hello".to_vec()); + set_disks + .put_object(bucket, "object", &mut reader, &ObjectOptions::default()) + .await + .expect("object should be written"); + + let list_result = set_disks + .clone() + .list_objects_v2(bucket, "", None, None, 1000, false, None, false) + .await + .expect("set-level list_objects_v2 should succeed"); + assert_eq!(list_result.objects.len(), 1); + assert_eq!(list_result.objects[0].name, "object"); + + let versions_result = set_disks + .clone() + .list_object_versions(bucket, "", None, None, None, 1000) + .await + .expect("set-level list_object_versions should succeed"); + assert_eq!(versions_result.objects.len(), 1); + assert_eq!(versions_result.objects[0].name, "object"); + + let (tx, mut rx) = mpsc::channel(4); + set_disks + .clone() + .walk(CancellationToken::new(), bucket, "", tx, WalkOptions::default()) + .await + .expect("set-level walk should succeed"); + + let mut walked_names = Vec::new(); + while let Some(item) = rx.recv().await { + if let Some(object) = item.item { + walked_names.push(object.name); + } + } + assert!(walked_names.iter().any(|name| name == "object")); + } + + #[tokio::test] + async fn set_level_heal_format_repairs_unformatted_disk() { + let set_disks = make_local_bucket_test_set_disks_with_missing_format().await; + let disk = { + let disks = set_disks.disks.read().await; + disks[1].clone().expect("second disk should exist") + }; + + let before = load_format_erasure(&disk, true) + .await + .expect_err("second disk should start unformatted"); + assert_eq!(before, DiskError::UnformattedDisk); + + let (heal_result, heal_err) = set_disks.heal_format(false).await.expect("heal_format should complete"); + assert!(heal_err.is_none(), "heal_format should repair the local unformatted disk"); + assert_eq!(heal_result.disk_count, 2); + assert_eq!(heal_result.set_count, 1); + assert_eq!(heal_result.after.drives[1].state, DriveState::Ok.to_string()); + + let repaired = load_format_erasure(&disk, true) + .await + .expect("second disk should contain a healed format"); + assert_eq!(repaired.erasure.this, set_disks.format.erasure.sets[0][1]); + } + + #[tokio::test] + async fn remaining_unsupported_trait_stubs_return_typed_errors() { + let set_disks = make_test_set_disks(Vec::new()).await; + + let (heal_result, heal_err) = make_local_bucket_test_set_disks() + .await + .heal_format(false) + .await + .expect("heal_format should be callable on formatted disks"); + assert!(matches!(heal_err, Some(StorageError::NoHealRequired))); + assert_eq!(heal_result.disk_count, 2); + + let copy_part_err = set_disks + .copy_object_part( + "bucket", + "src", + "bucket", + "dst", + "upload-id", + 1, + 0, + 1, + &ObjectInfo::default(), + &ObjectOptions::default(), + &ObjectOptions::default(), + ) + .await + .expect_err("unsupported copy_object_part should return a typed error"); + assert!(matches!(copy_part_err, StorageError::NotImplemented)); + + let abandoned_err = set_disks + .check_abandoned_parts("bucket", "object", &HealOpts::default()) + .await + .expect_err("abandoned-parts check should stay in the upper reconciliation layer"); + assert!(matches!(abandoned_err, StorageError::NotImplemented)); + } } diff --git a/crates/ecstore/src/sets.rs b/crates/ecstore/src/sets.rs index 12b20e78b..e362ddd52 100644 --- a/crates/ecstore/src/sets.rs +++ b/crates/ecstore/src/sets.rs @@ -440,22 +440,62 @@ impl BucketOperations for Sets { type Error = Error; #[tracing::instrument(skip(self))] - async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> { - unimplemented!() + async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> { + for set in &self.disk_set { + set.make_bucket(bucket, opts).await?; + } + + Ok(()) } #[tracing::instrument(skip(self))] - async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { - unimplemented!() + async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result { + let mut first_err = None; + for set in &self.disk_set { + match set.get_bucket_info(bucket, opts).await { + Ok(info) => return Ok(info), + Err(err) if first_err.is_none() => first_err = Some(err), + Err(_) => {} + } + } + + Err(first_err.unwrap_or_else(|| StorageError::BucketNotFound(bucket.to_string()))) } #[tracing::instrument(skip(self))] - async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { - unimplemented!() + async fn list_bucket(&self, opts: &BucketOptions) -> Result> { + let mut buckets = HashMap::new(); + let mut first_err = None; + + for set in &self.disk_set { + match set.list_bucket(opts).await { + Ok(set_buckets) => { + for bucket in set_buckets { + buckets.entry(bucket.name.clone()).or_insert(bucket); + } + } + Err(err) if first_err.is_none() => first_err = Some(err), + Err(_) => {} + } + } + + if buckets.is_empty() + && let Some(err) = first_err + { + return Err(err); + } + + let mut buckets = buckets.into_values().collect::>(); + buckets.sort_by(|left, right| left.name.cmp(&right.name)); + Ok(buckets) } #[tracing::instrument(skip(self))] - async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { - unimplemented!() + async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { + for set in &self.disk_set { + set.delete_bucket(bucket, opts).await?; + } + + Ok(()) } } @@ -543,8 +583,10 @@ impl rustfs_storage_api::ObjectOperations for Sets { } #[tracing::instrument(skip(self))] - async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, _force_del_marker: bool) -> Result<()> { - unimplemented!() + async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> { + self.get_disks_by_key(object) + .delete_object_version(bucket, object, fi, force_del_marker) + .await } #[tracing::instrument(skip(self))] @@ -678,40 +720,51 @@ impl rustfs_storage_api::ListOperations for Sets { #[tracing::instrument(skip(self))] async fn list_objects_v2( self: Arc, - _bucket: &str, - _prefix: &str, - _continuation_token: Option, - _delimiter: Option, - _max_keys: i32, - _fetch_owner: bool, - _start_after: Option, - _incl_deleted: bool, + bucket: &str, + prefix: &str, + continuation_token: Option, + delimiter: Option, + max_keys: i32, + fetch_owner: bool, + start_after: Option, + incl_deleted: bool, ) -> Result { - unimplemented!() + self.inner_list_objects_v2( + bucket, + prefix, + continuation_token, + delimiter, + max_keys, + fetch_owner, + start_after, + incl_deleted, + ) + .await } #[tracing::instrument(skip(self))] async fn list_object_versions( self: Arc, - _bucket: &str, - _prefix: &str, - _marker: Option, - _version_marker: Option, - _delimiter: Option, - _max_keys: i32, + bucket: &str, + prefix: &str, + marker: Option, + version_marker: Option, + delimiter: Option, + max_keys: i32, ) -> Result { - unimplemented!() + self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys) + .await } async fn walk( self: Arc, - _rx: CancellationToken, - _bucket: &str, - _prefix: &str, - _result: tokio::sync::mpsc::Sender, - _opts: WalkOptions, + rx: CancellationToken, + bucket: &str, + prefix: &str, + result: tokio::sync::mpsc::Sender, + opts: WalkOptions, ) -> Result<()> { - unimplemented!() + self.walk_internal(rx, bucket, prefix, result, opts).await } } @@ -762,7 +815,7 @@ impl rustfs_storage_api::MultipartOperations for Sets { _src_opts: &ObjectOptions, _dst_opts: &ObjectOptions, ) -> Result<()> { - unimplemented!() + Err(StorageError::NotImplemented) } #[tracing::instrument(skip(self))] @@ -920,8 +973,22 @@ impl rustfs_storage_api::HealOperations for Sets { Ok((res, None)) } #[tracing::instrument(skip(self))] - async fn heal_bucket(&self, _bucket: &str, _opts: &HealOpts) -> Result { - unimplemented!() + async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result { + let mut result = HealResultItem { + heal_item_type: HealItemType::Bucket.to_string(), + bucket: bucket.to_string(), + set_count: self.set_count, + ..Default::default() + }; + + for set in &self.disk_set { + let mut set_result = set.heal_bucket(bucket, opts).await?; + result.disk_count += set_result.disk_count; + result.before.drives.append(&mut set_result.before.drives); + result.after.drives.append(&mut set_result.after.drives); + } + + Ok(result) } #[tracing::instrument(skip(self))] async fn heal_object( @@ -936,12 +1003,22 @@ impl rustfs_storage_api::HealOperations for Sets { .await } #[tracing::instrument(skip(self))] - async fn get_pool_and_set(&self, _id: &str) -> Result<(Option, Option, Option)> { - unimplemented!() + async fn get_pool_and_set(&self, id: &str) -> Result<(Option, Option, Option)> { + for (set_idx, set) in self.format.erasure.sets.iter().enumerate() { + for (disk_idx, disk_id) in set.iter().enumerate() { + if disk_id.to_string() == id { + return Ok((Some(self.pool_idx), Some(set_idx), Some(disk_idx))); + } + } + } + + Err(Error::DiskNotFound) } #[tracing::instrument(skip(self))] async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> { - unimplemented!() + // Multipart orphan reconciliation is intentionally retained above the pool/set layers + // until there is a concrete caller and a stable lower-level contract to implement. + Err(StorageError::NotImplemented) } } @@ -1021,6 +1098,9 @@ async fn init_storage_disks_with_errors( #[cfg(test)] mod tests { use super::*; + use crate::layout::endpoint::Endpoint; + use rustfs_storage_api::HealOperations as _; + use rustfs_storage_api::ListOperations as _; #[test] fn test_apply_delete_objects_results_preserves_original_order_for_out_of_order_batches() { @@ -1089,4 +1169,164 @@ mod tests { Some(Error::other("third failed").to_string()) ); } + + #[tokio::test] + async fn sets_get_pool_and_set_returns_matching_coordinates() { + let format = FormatV3::new(2, 2); + let target = format.erasure.sets[1][0].to_string(); + + let endpoints = vec![ + Endpoint::try_from("http://127.0.0.1:9000/data0").expect("first endpoint should parse"), + Endpoint::try_from("http://127.0.0.1:9001/data1").expect("second endpoint should parse"), + Endpoint::try_from("http://127.0.0.1:9002/data2").expect("third endpoint should parse"), + Endpoint::try_from("http://127.0.0.1:9003/data3").expect("fourth endpoint should parse"), + ]; + + let sets = Sets { + id: format.id, + disk_set: Vec::new(), + pool_idx: 3, + endpoints: PoolEndpoints { + legacy: false, + set_count: 2, + drives_per_set: 2, + endpoints: Endpoints::from(endpoints), + cmd_line: String::new(), + platform: String::new(), + }, + format, + parity_count: 1, + set_count: 2, + set_drive_count: 2, + default_parity_count: 1, + distribution_algo: DistributionAlgoVersion::V1, + exit_signal: None, + }; + + let result = sets + .get_pool_and_set(&target) + .await + .expect("disk id should resolve within the pool"); + + assert_eq!(result, (Some(3), Some(1), Some(0))); + } + + #[tokio::test] + async fn sets_list_objects_v2_lists_objects_within_the_pool() { + let format = FormatV3::new(1, 2); + let mut endpoints = Vec::new(); + let mut disks = Vec::new(); + + for disk_idx in 0..2 { + let dir = tempfile::tempdir().expect("tempdir should be created"); + let mut endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(disk_idx); + + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("disk should be created"); + + let mut disk_format = format.clone(); + disk_format.erasure.this = format.erasure.sets[0][disk_idx]; + save_format_file(&Some(disk.clone()), &Some(disk_format)) + .await + .expect("format should be saved"); + + std::mem::forget(dir); + endpoints.push(endpoint); + disks.push(Some(disk)); + } + + let set_disks = SetDisks::new( + "test-owner".to_string(), + Arc::new(RwLock::new(disks)), + 2, + 1, + 0, + 0, + endpoints.clone(), + format.clone(), + Vec::new(), + ) + .await; + + let sets = Arc::new(Sets { + id: format.id, + disk_set: vec![set_disks], + pool_idx: 0, + endpoints: PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 2, + endpoints: Endpoints::from(endpoints), + cmd_line: String::new(), + platform: String::new(), + }, + format, + parity_count: 1, + set_count: 1, + set_drive_count: 2, + default_parity_count: 1, + distribution_algo: DistributionAlgoVersion::V1, + exit_signal: None, + }); + + sets.make_bucket("bucket", &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + + let mut reader = PutObjReader::from_vec(b"hello".to_vec()); + sets.put_object("bucket", "object", &mut reader, &ObjectOptions::default()) + .await + .expect("object should be written"); + + let result = sets + .clone() + .list_objects_v2("bucket", "", None, None, 1000, false, None, false) + .await + .expect("pool-level listing should succeed"); + + assert_eq!(result.objects.len(), 1); + assert_eq!(result.objects[0].name, "object"); + } + + #[tokio::test] + async fn sets_check_abandoned_parts_returns_typed_not_implemented_error() { + let format = FormatV3::new(1, 1); + let sets = Sets { + id: format.id, + disk_set: Vec::new(), + pool_idx: 0, + endpoints: PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 1, + endpoints: Endpoints::from(Vec::new()), + cmd_line: String::new(), + platform: String::new(), + }, + format, + parity_count: 0, + set_count: 1, + set_drive_count: 1, + default_parity_count: 0, + distribution_algo: DistributionAlgoVersion::V1, + exit_signal: None, + }; + + let err = sets + .check_abandoned_parts("bucket", "object", &HealOpts::default()) + .await + .expect_err("abandoned-parts ownership should stay above the pool/set storage layers"); + assert!(matches!(err, StorageError::NotImplemented)); + } } diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 2e698c514..7fa917050 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -156,23 +156,12 @@ impl ECStore { #[instrument(skip(self))] pub(super) async fn handle_check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> { - let object = encode_dir_object(object); - if self.single_pool() { - return self.pools[0].check_abandoned_parts(bucket, &object, opts).await; - } - - let mut errs = Vec::new(); - for pool in self.pools.iter() { - //TODO: IsSuspended - if let Err(err) = pool.check_abandoned_parts(bucket, &object, opts).await { - errs.push(err); - } - } - - if !errs.is_empty() { - return Err(errs[0].clone()); - } - - Ok(()) + let _ = (bucket, object, opts); + // Stale multipart reconciliation is already owned by the lifecycle-driven + // background cleanup path in `bucket_lifecycle_ops.rs`. There is currently + // no stable object-heal contract that should fan this request out through + // pool/set storage layers, so keep the placeholder explicit at the ECStore + // boundary instead of dispatching into lower layers. + Err(StorageError::NotImplemented) } } diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index a39a8324b..361d99108 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -180,9 +180,8 @@ impl ECStore { ) -> Result<()> { check_new_multipart_args(src_bucket, src_object)?; - // TODO: PutObjectReader - // self.put_object_part(dst_bucket, dst_object, upload_id, part_id, data, opts) - + // The full UploadPartCopy path still requires the higher S3/request layer to + // derive encryption, compression, and multipart checksum write semantics. Err(StorageError::NotImplemented) } diff --git a/crates/ecstore/src/store_list_objects.rs b/crates/ecstore/src/store_list_objects.rs index 806b0f6fe..b67beb3ac 100644 --- a/crates/ecstore/src/store_list_objects.rs +++ b/crates/ecstore/src/store_list_objects.rs @@ -23,6 +23,7 @@ use crate::error::{ }; use crate::object_api::{ObjectInfo, ObjectOptions}; use crate::set_disk::SetDisks; +use crate::sets::Sets; use crate::store::ECStore; use crate::store_utils::is_reserved_or_invalid_bucket; use futures::future::join_all; @@ -1482,7 +1483,1308 @@ async fn merge_entry_channels( } } +impl Sets { + #[allow(clippy::too_many_arguments)] + pub async fn inner_list_objects_v2( + self: Arc, + bucket: &str, + prefix: &str, + continuation_token: Option, + delimiter: Option, + max_keys: i32, + _fetch_owner: bool, + start_after: Option, + incl_deleted: bool, + ) -> Result { + let marker = if continuation_token.is_none() { + start_after + } else { + continuation_token.clone() + }; + + let loi = self + .list_objects_generic(bucket, prefix, marker, delimiter, max_keys, incl_deleted) + .await?; + Ok(ListObjectsV2Info { + is_truncated: loi.is_truncated, + continuation_token, + next_continuation_token: loi.next_marker, + objects: loi.objects, + prefixes: loi.prefixes, + }) + } + + pub async fn list_objects_generic( + self: Arc, + bucket: &str, + prefix: &str, + marker: Option, + delimiter: Option, + max_keys: i32, + incl_deleted: bool, + ) -> Result { + let max_keys = normalize_max_keys(max_keys); + let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) }; + let opts = ListPathOptions { + bucket: bucket.to_owned(), + prefix: prefix.to_owned(), + separator: delimiter.clone(), + limit: effective_max_keys, + marker, + incl_deleted, + ask_disks: list_quorum_from_env(), + ..Default::default() + }; + + if !opts.prefix.is_empty() && max_keys == 1 && opts.marker.is_none() { + match self + .get_object_info( + &opts.bucket, + &opts.prefix, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(res) => { + return Ok(ListObjectsInfo { + objects: vec![res], + ..Default::default() + }); + } + Err(err) => { + if is_err_bucket_not_found(&err) { + return Err(err); + } + } + }; + } + + let mut list_result = self + .list_path(&opts) + .await + .unwrap_or_else(|err| MetaCacheEntriesSortedResult { + err: Some(err.into()), + ..Default::default() + }); + + let disk_has_more = list_result.err.is_none(); + + if let Some(err) = list_result.err.take() + && err != rustfs_filemeta::Error::Unexpected + { + return Err(to_object_err(err.into(), vec![bucket, prefix])); + } + + if let Some(result) = list_result.entries.as_mut() { + result.forward_past(opts.marker); + } + + let mut get_objects = ObjectInfo::from_meta_cache_entries_sorted_infos( + &list_result.entries.unwrap_or_default(), + bucket, + prefix, + delimiter.clone(), + ) + .await; + + let mut is_truncated = false; + if max_keys <= 0 { + get_objects.clear(); + } else if get_objects.len() > max_keys as usize { + is_truncated = true; + get_objects.truncate(max_keys as usize); + } + + let mut next_marker = if is_truncated { + get_objects.last().map(|last| last.name.clone()) + } else { + None + }; + + let mut prefixes: Vec = Vec::new(); + let mut prefix_set: HashSet = HashSet::new(); + let mut objects = Vec::with_capacity(get_objects.len()); + for obj in get_objects { + if delimiter.is_some() { + if obj.is_dir && obj.mod_time.is_none() { + if prefix_set.insert(obj.name.clone()) { + prefixes.push(obj.name); + } + } else { + objects.push(obj); + } + } else { + objects.push(obj); + } + } + + if !is_truncated && disk_has_more { + let visible_count = objects.len() + prefixes.len(); + let should_truncate = if delimiter.is_none() { + visible_count > 0 + } else { + visible_count >= max_keys as usize + }; + if should_truncate { + is_truncated = true; + next_marker = objects + .last() + .map(|last| last.name.clone()) + .or_else(|| prefixes.last().cloned()); + } + } + + Ok(ListObjectsInfo { + is_truncated, + next_marker, + objects, + prefixes, + }) + } + + pub async fn inner_list_object_versions( + self: Arc, + bucket: &str, + prefix: &str, + marker: Option, + version_marker: Option, + delimiter: Option, + max_keys: i32, + ) -> Result { + let max_keys = normalize_max_keys(max_keys); + if marker.is_none() && version_marker.is_some() { + return Err(StorageError::NotImplemented); + } + + let has_version_marker = version_marker.is_some(); + let version_marker = if let Some(marker) = version_marker { + Some(parse_version_marker(marker)?) + } else { + None + }; + + let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) }; + let opts = ListPathOptions { + bucket: bucket.to_owned(), + prefix: prefix.to_owned(), + separator: delimiter.clone(), + limit: effective_max_keys, + marker, + incl_deleted: true, + ask_disks: list_quorum_from_env(), + versioned: true, + include_marker: has_version_marker, + ..Default::default() + }; + + let mut list_result = self + .list_path(&opts) + .await + .unwrap_or_else(|err| MetaCacheEntriesSortedResult { + err: Some(err.into()), + ..Default::default() + }); + + let disk_has_more = list_result.err.is_none(); + + if let Some(err) = list_result.err.take() + && err != rustfs_filemeta::Error::Unexpected + { + return Err(to_object_err(err.into(), vec![bucket, prefix])); + } + + if let Some(result) = list_result.entries.as_mut() + && !has_version_marker + { + result.forward_past(opts.marker.clone()); + } + + let version_marker = version_marker_for_entries(list_result.entries.as_ref(), opts.marker.as_deref(), version_marker); + + let mut get_objects = ObjectInfo::from_meta_cache_entries_sorted_versions( + &list_result.entries.unwrap_or_default(), + bucket, + prefix, + delimiter.clone(), + version_marker, + ) + .await; + + let mut is_truncated = false; + if max_keys <= 0 { + get_objects.clear(); + } else if get_objects.len() > max_keys as usize { + is_truncated = true; + get_objects.truncate(max_keys as usize); + } + + let mut next_marker: Option = None; + let mut next_version_idmarker: Option = None; + if is_truncated && let Some(last) = get_objects.last() { + next_marker = Some(last.name.clone()); + next_version_idmarker = Some(last.version_id.map(|v| v.to_string()).unwrap_or_else(|| "null".to_string())); + } + + let mut prefixes: Vec = Vec::new(); + let mut prefix_set: HashSet = HashSet::new(); + let mut objects = Vec::with_capacity(get_objects.len()); + for obj in get_objects { + if delimiter.is_some() { + if obj.is_dir && obj.mod_time.is_none() { + if prefix_set.insert(obj.name.clone()) { + prefixes.push(obj.name); + } + } else { + objects.push(obj); + } + } else { + objects.push(obj); + } + } + + if !is_truncated && disk_has_more { + let visible_count = objects.len() + prefixes.len(); + let should_truncate = if delimiter.is_none() { + visible_count > 0 + } else { + visible_count >= max_keys as usize + }; + if should_truncate { + is_truncated = true; + if let Some(last) = objects.last() { + next_marker = Some(last.name.clone()); + next_version_idmarker = Some(last.version_id.map(|v| v.to_string()).unwrap_or_else(|| "null".to_string())); + } else if let Some(last_prefix) = prefixes.last().cloned() { + next_marker = Some(last_prefix); + next_version_idmarker = None; + } + } + } + + Ok(ListObjectVersionsInfo { + is_truncated, + next_marker, + next_version_idmarker, + objects, + prefixes, + }) + } + + pub async fn list_path(self: Arc, o: &ListPathOptions) -> Result { + check_list_objs_args(&o.bucket, &o.prefix, &o.marker)?; + + let mut o = o.clone(); + o.marker = o.marker.filter(|v| v >= &o.prefix); + + if let Some(marker) = &o.marker + && !o.prefix.is_empty() + && !marker.starts_with(&o.prefix) + { + return Err(Error::Unexpected); + } + + if o.limit == 0 { + return Err(Error::Unexpected); + } + + if o.prefix.starts_with(SLASH_SEPARATOR) { + return Err(Error::Unexpected); + } + + let slash_separator = Some(SLASH_SEPARATOR.to_owned()); + o.include_directories = o.separator == slash_separator; + + if (o.separator == slash_separator || o.separator.is_none()) && !o.recursive { + o.recursive = o.separator != slash_separator; + o.separator = slash_separator; + } else { + o.recursive = true; + } + + o.parse_marker(); + + if o.base_dir.is_empty() { + o.base_dir = base_dir_from_prefix(&o.prefix); + } + + o.transient = o.transient || is_reserved_or_invalid_bucket(&o.bucket, false); + o.set_filter(); + if o.transient { + o.create = false; + } + + let cancel = CancellationToken::new(); + let (err_tx, mut err_rx) = broadcast::channel::>(1); + let (sender, recv) = mpsc::channel(o.limit as usize); + + let sets = self.clone(); + let opts = o.clone(); + let cancel_rx1 = cancel.clone(); + let cancel_rx1_for_err = cancel_rx1.clone(); + let err_tx1 = err_tx.clone(); + let job1 = tokio::spawn( + async move { + let mut opts = opts; + opts.stop_disk_at_limit = true; + if let Err(err) = sets.list_merged(cancel_rx1, opts, sender).await + && !cancel_rx1_for_err.is_cancelled() + { + error!("list_merged err {:?}", err); + let _ = err_tx1.send(Arc::new(err)); + } + } + .instrument(tracing::Span::current()), + ); + + let cancel_rx2 = cancel.clone(); + let (result_tx, mut result_rx) = mpsc::channel(1); + let err_tx2 = err_tx.clone(); + let opts = o.clone(); + let job2 = tokio::spawn( + async move { + if let Err(err) = gather_results(cancel_rx2, opts, recv, result_tx).await { + error!("gather_results err {:?}", err); + let _ = err_tx2.send(Arc::new(err)); + } + cancel.cancel(); + } + .instrument(tracing::Span::current()), + ); + + let mut result = tokio::select! { + res = err_rx.recv() => { + match res { + Ok(err) => MetaCacheEntriesSortedResult { entries: None, err: Some(err.as_ref().clone().into()) }, + Err(err) => MetaCacheEntriesSortedResult { entries: None, err: Some(rustfs_filemeta::Error::other(err)) }, + } + } + Some(result) = result_rx.recv() => result, + }; + + join_all(vec![job1, job2]).await; + + if let Ok(err) = err_rx.try_recv() { + result.err = Some(err.as_ref().clone().into()); + } + + if result.err.is_some() { + return Ok(result); + } + + if let Some(entries) = result.entries.as_mut() { + entries.reuse = true; + let truncated = !entries.entries().is_empty() || result.err.is_none(); + entries.o.0.truncate(o.limit as usize); + if !o.transient && truncated { + entries.list_id = if let Some(id) = o.id { + Some(id) + } else { + Some(Uuid::new_v4().to_string()) + }; + } + + if !truncated { + result.err = Some(Error::Unexpected.into()); + } + } + + Ok(result) + } + + async fn list_merged( + &self, + rx: CancellationToken, + opts: ListPathOptions, + sender: Sender, + ) -> Result> { + let mut futures = Vec::new(); + let mut inputs = Vec::new(); + + for set in &self.disk_set { + let (send, recv) = mpsc::channel(100); + inputs.push(recv); + let opts = opts.clone(); + let rx_clone = rx.clone(); + let set = set.clone(); + futures.push(async move { set.list_path(rx_clone, opts, send).await }); + } + + tokio::spawn( + async move { + if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await { + error!("merge_entry_channels err {:?}", err); + } + } + .instrument(tracing::Span::current()), + ); + + let results = join_all(futures).await; + let mut all_at_eof = true; + let mut errs = Vec::new(); + for result in results { + if let Err(err) = result { + all_at_eof = false; + errs.push(Some(err)); + } else { + errs.push(None); + } + } + + if is_all_not_found(&errs) { + return Ok(Vec::new()); + } + + for err in &errs { + if let Some(err) = err { + if err == &Error::Unexpected { + continue; + } + return Err(err.clone()); + } else { + all_at_eof = false; + } + } + + _ = all_at_eof; + Ok(Vec::new()) + } + + #[allow(unused_assignments)] + pub async fn walk_internal( + self: Arc, + rx: CancellationToken, + 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 set in &self.disk_set { + let (mut disks, infos, _) = set.get_online_disks_with_healing_and_info(true).await; + let opts = opts.clone(); + let (sender, list_out_rx) = mpsc::channel::(1); + inputs.push(list_out_rx); + let rx_clone = rx.clone(); + let set = set.clone(); + 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().div_ceil(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 = rand::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); + ensure_non_empty_listing_disks(bucket, &path, &disks)?; + + 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 tx1 = sender.clone(); + let tx2 = sender.clone(); + + list_path_raw( + rx_clone, + ListPathRawOptions { + disks: disks.iter().cloned().map(Some).collect(), + fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), + bucket: bucket.to_owned(), + path, + recursive: 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 Some(entry) = entries.resolve(resolver) + && 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 bucket_clone = bucket.clone(); + + 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_clone) { + Ok(res) => res, + Err(err) => { + if !sent_err { + let item = ObjectInfoOrErr { + item: None, + err: Some(err.into()), + }; + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + sent_err = true; + return; + } + continue; + } + }; + + if let Some(filter) = opts.filter { + if filter(&fi) { + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &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(ObjectInfo::from_file_info(&fi, &bucket_clone, &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 if opts.include_free_versions { + entry.file_info_versions_with_free_versions(&bucket_clone) + } else { + entry.file_info_versions(&bucket_clone) + } { + Ok(res) => res, + Err(err) => { + let item = ObjectInfoOrErr { + item: None, + err: Some(err.into()), + }; + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + return; + } + }; + + for fi in &fvs.versions { + if let Some(filter) = opts.filter { + if filter(fi) { + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &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(ObjectInfo::from_file_info(fi, &bucket_clone, &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); + } + } + } + + if opts.include_free_versions { + for fi in &fvs.free_versions { + if let Some(filter) = opts.filter { + if filter(fi) { + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &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(ObjectInfo::from_file_info(fi, &bucket_clone, &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); + } + } + } + } + } + } + .instrument(tracing::Span::current()), + ); + + tokio::spawn(async move { merge_entry_channels(rx, inputs, merge_tx, 1).await }.instrument(tracing::Span::current())); + + let walk_started = std::time::Instant::now(); + let walk_results = join_all(futures).await; + let mut errs = Vec::new(); + for walk_result in walk_results { + match walk_result { + Ok(()) => errs.push(None), + Err(err) => errs.push(Some(err.into())), + } + } + rustfs_io_metrics::record_stage_duration( + "sets_list_objects_walk_internal", + walk_started.elapsed().as_secs_f64() * 1000.0, + ); + + let result = walk_result_from_set_errors(&errs); + if let Err(err) = &result { + error!( + bucket = %bucket, + prefix = %prefix, + error = ?err, + set_errors = ?errs, + "walk_internal list_path_raw tasks failed" + ); + } + + result + } +} + impl SetDisks { + #[allow(clippy::too_many_arguments)] + pub async fn inner_list_objects_v2( + self: Arc, + bucket: &str, + prefix: &str, + continuation_token: Option, + delimiter: Option, + max_keys: i32, + _fetch_owner: bool, + start_after: Option, + incl_deleted: bool, + ) -> Result { + let marker = if continuation_token.is_none() { + start_after + } else { + continuation_token.clone() + }; + + let loi = self + .list_objects_generic(bucket, prefix, marker, delimiter, max_keys, incl_deleted) + .await?; + Ok(ListObjectsV2Info { + is_truncated: loi.is_truncated, + continuation_token, + next_continuation_token: loi.next_marker, + objects: loi.objects, + prefixes: loi.prefixes, + }) + } + + pub async fn list_objects_generic( + self: Arc, + bucket: &str, + prefix: &str, + marker: Option, + delimiter: Option, + max_keys: i32, + incl_deleted: bool, + ) -> Result { + let max_keys = normalize_max_keys(max_keys); + let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) }; + let opts = ListPathOptions { + bucket: bucket.to_owned(), + prefix: prefix.to_owned(), + separator: delimiter.clone(), + limit: effective_max_keys, + marker, + incl_deleted, + ask_disks: list_quorum_from_env(), + ..Default::default() + }; + + if !opts.prefix.is_empty() && max_keys == 1 && opts.marker.is_none() { + match self + .get_object_info( + &opts.bucket, + &opts.prefix, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + { + Ok(res) => { + return Ok(ListObjectsInfo { + objects: vec![res], + ..Default::default() + }); + } + Err(err) => { + if is_err_bucket_not_found(&err) { + return Err(err); + } + } + }; + } + + let mut list_result = self + .list_path_result(&opts) + .await + .unwrap_or_else(|err| MetaCacheEntriesSortedResult { + err: Some(err.into()), + ..Default::default() + }); + + let disk_has_more = list_result.err.is_none(); + + if let Some(err) = list_result.err.take() + && err != rustfs_filemeta::Error::Unexpected + { + return Err(to_object_err(err.into(), vec![bucket, prefix])); + } + + if let Some(result) = list_result.entries.as_mut() { + result.forward_past(opts.marker); + } + + let mut get_objects = ObjectInfo::from_meta_cache_entries_sorted_infos( + &list_result.entries.unwrap_or_default(), + bucket, + prefix, + delimiter.clone(), + ) + .await; + + let mut is_truncated = false; + if max_keys <= 0 { + get_objects.clear(); + } else if get_objects.len() > max_keys as usize { + is_truncated = true; + get_objects.truncate(max_keys as usize); + } + + let mut next_marker = if is_truncated { + get_objects.last().map(|last| last.name.clone()) + } else { + None + }; + + let mut prefixes: Vec = Vec::new(); + let mut prefix_set: HashSet = HashSet::new(); + let mut objects = Vec::with_capacity(get_objects.len()); + for obj in get_objects { + if delimiter.is_some() { + if obj.is_dir && obj.mod_time.is_none() { + if prefix_set.insert(obj.name.clone()) { + prefixes.push(obj.name); + } + } else { + objects.push(obj); + } + } else { + objects.push(obj); + } + } + + if !is_truncated && disk_has_more { + let visible_count = objects.len() + prefixes.len(); + let should_truncate = if delimiter.is_none() { + visible_count > 0 + } else { + visible_count >= max_keys as usize + }; + if should_truncate { + is_truncated = true; + next_marker = objects + .last() + .map(|last| last.name.clone()) + .or_else(|| prefixes.last().cloned()); + } + } + + Ok(ListObjectsInfo { + is_truncated, + next_marker, + objects, + prefixes, + }) + } + + pub async fn inner_list_object_versions( + self: Arc, + bucket: &str, + prefix: &str, + marker: Option, + version_marker: Option, + delimiter: Option, + max_keys: i32, + ) -> Result { + let max_keys = normalize_max_keys(max_keys); + if marker.is_none() && version_marker.is_some() { + return Err(StorageError::NotImplemented); + } + + let has_version_marker = version_marker.is_some(); + let version_marker = if let Some(marker) = version_marker { + Some(parse_version_marker(marker)?) + } else { + None + }; + + let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) }; + let opts = ListPathOptions { + bucket: bucket.to_owned(), + prefix: prefix.to_owned(), + separator: delimiter.clone(), + limit: effective_max_keys, + marker, + incl_deleted: true, + ask_disks: list_quorum_from_env(), + versioned: true, + include_marker: has_version_marker, + ..Default::default() + }; + + let mut list_result = self + .list_path_result(&opts) + .await + .unwrap_or_else(|err| MetaCacheEntriesSortedResult { + err: Some(err.into()), + ..Default::default() + }); + + let disk_has_more = list_result.err.is_none(); + + if let Some(err) = list_result.err.take() + && err != rustfs_filemeta::Error::Unexpected + { + return Err(to_object_err(err.into(), vec![bucket, prefix])); + } + + if let Some(result) = list_result.entries.as_mut() + && !has_version_marker + { + result.forward_past(opts.marker.clone()); + } + + let version_marker = version_marker_for_entries(list_result.entries.as_ref(), opts.marker.as_deref(), version_marker); + + let mut get_objects = ObjectInfo::from_meta_cache_entries_sorted_versions( + &list_result.entries.unwrap_or_default(), + bucket, + prefix, + delimiter.clone(), + version_marker, + ) + .await; + + let mut is_truncated = false; + if max_keys <= 0 { + get_objects.clear(); + } else if get_objects.len() > max_keys as usize { + is_truncated = true; + get_objects.truncate(max_keys as usize); + } + + let mut next_marker: Option = None; + let mut next_version_idmarker: Option = None; + if is_truncated && let Some(last) = get_objects.last() { + next_marker = Some(last.name.clone()); + next_version_idmarker = Some(last.version_id.map(|v| v.to_string()).unwrap_or_else(|| "null".to_string())); + } + + let mut prefixes: Vec = Vec::new(); + let mut prefix_set: HashSet = HashSet::new(); + let mut objects = Vec::with_capacity(get_objects.len()); + for obj in get_objects { + if delimiter.is_some() { + if obj.is_dir && obj.mod_time.is_none() { + if prefix_set.insert(obj.name.clone()) { + prefixes.push(obj.name); + } + } else { + objects.push(obj); + } + } else { + objects.push(obj); + } + } + + if !is_truncated && disk_has_more { + let visible_count = objects.len() + prefixes.len(); + let should_truncate = if delimiter.is_none() { + visible_count > 0 + } else { + visible_count >= max_keys as usize + }; + if should_truncate { + is_truncated = true; + if let Some(last) = objects.last() { + next_marker = Some(last.name.clone()); + next_version_idmarker = Some(last.version_id.map(|v| v.to_string()).unwrap_or_else(|| "null".to_string())); + } else if let Some(last_prefix) = prefixes.last().cloned() { + next_marker = Some(last_prefix); + next_version_idmarker = None; + } + } + } + + Ok(ListObjectVersionsInfo { + is_truncated, + next_marker, + next_version_idmarker, + objects, + prefixes, + }) + } + + pub async fn walk_internal( + self: Arc, + rx: CancellationToken, + bucket: &str, + prefix: &str, + result: Sender, + opts: WalkOptions, + ) -> Result<()> { + check_list_objs_args(bucket, prefix, &None)?; + + let (entry_tx, mut entry_rx) = mpsc::channel::(100); + let bucket_name = bucket.to_owned(); + let bucket_name_for_list = bucket_name.clone(); + + let versioning_config = match get_versioning_config(&bucket_name).await { + Ok((res, _)) => Some(res), + Err(_) => None, + }; + + let result_task = tokio::spawn( + async move { + while let Some(entry) = entry_rx.recv().await { + if opts.latest_only { + let fi = match entry.to_fileinfo(&bucket_name) { + Ok(res) => res, + Err(err) => { + let item = ObjectInfoOrErr { + item: None, + err: Some(err.into()), + }; + let _ = result.send(item).await; + return; + } + }; + + if let Some(filter) = opts.filter + && !filter(&fi) + { + continue; + } + + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(&fi, &bucket_name, &fi.name, { + if let Some(v) = &versioning_config { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + let _ = result.send(item).await; + continue; + } + + let versions = match if opts.include_free_versions { + entry.file_info_versions_with_free_versions(&bucket_name) + } else { + entry.file_info_versions(&bucket_name) + } { + Ok(res) => res, + Err(err) => { + let item = ObjectInfoOrErr { + item: None, + err: Some(err.into()), + }; + let _ = result.send(item).await; + return; + } + }; + + for fi in &versions.versions { + if let Some(filter) = opts.filter + && !filter(fi) + { + continue; + } + + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(fi, &bucket_name, &fi.name, { + if let Some(v) = &versioning_config { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + let _ = result.send(item).await; + } + + if opts.include_free_versions { + for fi in &versions.free_versions { + if let Some(filter) = opts.filter + && !filter(fi) + { + continue; + } + + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(fi, &bucket_name, &fi.name, { + if let Some(v) = &versioning_config { + v.versioned(&fi.name) + } else { + false + } + })), + err: None, + }; + let _ = result.send(item).await; + } + } + } + } + .instrument(tracing::Span::current()), + ); + + let limit = i32::try_from(opts.limit).unwrap_or(i32::MAX); + let list_result = self + .list_path( + rx, + ListPathOptions { + bucket: bucket_name_for_list, + prefix: prefix.to_owned(), + marker: opts.marker.clone(), + limit, + ask_disks: opts.ask_disks.clone(), + incl_deleted: true, + recursive: true, + versioned: true, + ..Default::default() + }, + entry_tx, + ) + .await; + + let _ = result_task.await; + + match list_result { + Ok(()) => Ok(()), + Err(err) => walk_result_from_set_errors(&[Some(err)]), + } + } + + pub async fn list_path_result(self: Arc, o: &ListPathOptions) -> Result { + check_list_objs_args(&o.bucket, &o.prefix, &o.marker)?; + + let mut o = o.clone(); + o.marker = o.marker.filter(|v| v >= &o.prefix); + + if let Some(marker) = &o.marker + && !o.prefix.is_empty() + && !marker.starts_with(&o.prefix) + { + return Err(Error::Unexpected); + } + + if o.limit == 0 { + return Err(Error::Unexpected); + } + + if o.prefix.starts_with(SLASH_SEPARATOR) { + return Err(Error::Unexpected); + } + + let slash_separator = Some(SLASH_SEPARATOR.to_owned()); + o.include_directories = o.separator == slash_separator; + + if (o.separator == slash_separator || o.separator.is_none()) && !o.recursive { + o.recursive = o.separator != slash_separator; + o.separator = slash_separator; + } else { + o.recursive = true; + } + + o.parse_marker(); + + if o.base_dir.is_empty() { + o.base_dir = base_dir_from_prefix(&o.prefix); + } + + o.transient = o.transient || is_reserved_or_invalid_bucket(&o.bucket, false); + o.set_filter(); + if o.transient { + o.create = false; + } + + let cancel = CancellationToken::new(); + let (err_tx, mut err_rx) = broadcast::channel::>(1); + let (sender, recv) = mpsc::channel(o.limit as usize); + + let set = self.clone(); + let opts = o.clone(); + let cancel_rx1 = cancel.clone(); + let cancel_rx1_for_err = cancel_rx1.clone(); + let err_tx1 = err_tx.clone(); + let job1 = tokio::spawn( + async move { + let mut opts = opts; + opts.stop_disk_at_limit = true; + if let Err(err) = set.list_path(cancel_rx1, opts, sender).await + && !cancel_rx1_for_err.is_cancelled() + { + error!("list_path err {:?}", err); + let _ = err_tx1.send(Arc::new(err)); + } + } + .instrument(tracing::Span::current()), + ); + + let cancel_rx2 = cancel.clone(); + let (result_tx, mut result_rx) = mpsc::channel(1); + let err_tx2 = err_tx.clone(); + let opts = o.clone(); + let job2 = tokio::spawn( + async move { + if let Err(err) = gather_results(cancel_rx2, opts, recv, result_tx).await { + error!("gather_results err {:?}", err); + let _ = err_tx2.send(Arc::new(err)); + } + cancel.cancel(); + } + .instrument(tracing::Span::current()), + ); + + let mut result = tokio::select! { + res = err_rx.recv() => { + match res { + Ok(err) => MetaCacheEntriesSortedResult { entries: None, err: Some(err.as_ref().clone().into()) }, + Err(err) => MetaCacheEntriesSortedResult { entries: None, err: Some(rustfs_filemeta::Error::other(err)) }, + } + } + Some(result) = result_rx.recv() => result, + }; + + join_all(vec![job1, job2]).await; + + if let Ok(err) = err_rx.try_recv() { + result.err = Some(err.as_ref().clone().into()); + } + + if result.err.is_some() { + return Ok(result); + } + + if let Some(entries) = result.entries.as_mut() { + entries.reuse = true; + let truncated = !entries.entries().is_empty() || result.err.is_none(); + entries.o.0.truncate(o.limit as usize); + if !o.transient && truncated { + entries.list_id = if let Some(id) = o.id { + Some(id) + } else { + Some(Uuid::new_v4().to_string()) + }; + } + + if !truncated { + result.err = Some(Error::Unexpected.into()); + } + } + + Ok(result) + } + pub async fn list_path(&self, rx: CancellationToken, opts: ListPathOptions, sender: Sender) -> Result<()> { let (mut disks, infos, _) = self.get_online_disks_with_healing_and_info(true).await; diff --git a/crates/iam/src/oidc.rs b/crates/iam/src/oidc.rs index d8e0306b4..69ada4694 100644 --- a/crates/iam/src/oidc.rs +++ b/crates/iam/src/oidc.rs @@ -1634,7 +1634,7 @@ mod tests { fn start_mock_oidc_discovery_server( build_discovery_issuer: F, max_requests: usize, - ) -> (String, std::thread::JoinHandle<()>) + ) -> Option<(String, std::thread::JoinHandle<()>)> where F: Fn(&str) -> String + Send + 'static, { @@ -1650,8 +1650,12 @@ mod tests { const IDLE_SHUTDOWN: Duration = Duration::from_secs(1); const ABSOLUTE_CAP: Duration = Duration::from_secs(5); - let listener = TcpListener::bind("127.0.0.1:0").unwrap(); - let base = format!("http://{}", listener.local_addr().unwrap()); + let listener = match TcpListener::bind("127.0.0.1:0") { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, + Err(err) => panic!("test listener should bind: {err}"), + }; + let base = format!("http://{}", listener.local_addr().expect("listener local address should be available")); let discovery_issuer = build_discovery_issuer(&base); let discovery_body = serde_json::json!({ "issuer": discovery_issuer, @@ -1754,7 +1758,7 @@ mod tests { .recv_timeout(Duration::from_millis(100)) .expect("mock OIDC discovery server should become ready"); - (base, handle) + Some((base, handle)) } fn discovery_error_contains_all_variants(err: &str, base: &str) -> bool { @@ -1776,7 +1780,9 @@ mod tests { async fn test_validate_oidc_provider_config_retries_with_issuer_candidates() { // Discovery document must advertise the canonical issuer path. The first candidate has no // trailing slash; openidconnect rejects issuer mismatch, then the second variant succeeds. - let (base, handle) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/rustfs/"), 8); + let Some((base, handle)) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/rustfs/"), 8) else { + return; + }; let config_url = format!("{base}/application/o/rustfs"); let config = build_mocked_oidc_provider_config("default", &config_url); @@ -1789,7 +1795,9 @@ mod tests { #[tokio::test] async fn test_validate_oidc_provider_config_returns_detailed_errors() { - let (base, handle) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/other"), 8); + let Some((base, handle)) = start_mock_oidc_discovery_server(|base| format!("{base}/application/o/other"), 8) else { + return; + }; let config_url = format!("{base}/application/o/rustfs"); let config = build_mocked_oidc_provider_config("default", &config_url); diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index ade420996..03b452e66 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -1016,9 +1016,13 @@ mod tests { StatusCode::OK } - async fn start_test_server(state: TestState) -> (String, tokio::task::JoinHandle<()>) { - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); + async fn start_test_server(state: TestState) -> Option<(String, tokio::task::JoinHandle<()>)> { + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); let app = Router::new() .route("/stream", get(get_stream).head(reject_head).put(accept_put)) .route("/stall", get(get_stalling_stream)) @@ -1028,7 +1032,7 @@ mod tests { axum::serve(listener, app).await.unwrap(); }); - (format!("http://{addr}/stream"), handle) + Some((format!("http://{addr}/stream"), handle)) } #[test] @@ -1055,7 +1059,9 @@ mod tests { #[tokio::test] async fn http_reader_does_not_send_preflight_head() { let state = TestState::default(); - let (url, handle) = start_test_server(state.clone()).await; + let Some((url, handle)) = start_test_server(state.clone()).await else { + return; + }; let mut reader = HttpReader::new(url, Method::GET, HeaderMap::new(), None).await.unwrap(); let mut buf = Vec::new(); @@ -1071,7 +1077,9 @@ mod tests { #[tokio::test] async fn http_reader_stall_timeout_triggers_after_progress_stops() { let state = TestState::default(); - let (base_url, handle) = start_test_server(state.clone()).await; + let Some((base_url, handle)) = start_test_server(state.clone()).await else { + return; + }; let url = base_url.replace("/stream", "/stall"); let mut reader = @@ -1099,7 +1107,9 @@ mod tests { #[tokio::test] async fn http_writer_does_not_send_empty_preflight_put() { let state = TestState::default(); - let (url, handle) = start_test_server(state.clone()).await; + let Some((url, handle)) = start_test_server(state.clone()).await else { + return; + }; let mut writer = HttpWriter::new(url, Method::PUT, HeaderMap::new()).await.unwrap(); writer.write_all(b"payload").await.unwrap(); @@ -1114,7 +1124,9 @@ mod tests { #[tokio::test] async fn http_writer_handles_many_small_writes() { let state = TestState::default(); - let (url, handle) = start_test_server(state.clone()).await; + let Some((url, handle)) = start_test_server(state.clone()).await else { + return; + }; let mut writer = HttpWriter::new(url, Method::PUT, HeaderMap::new()).await.unwrap(); let chunk = b"0123456789abcdef"; @@ -1134,7 +1146,9 @@ mod tests { #[tokio::test] async fn http_writer_supports_vectored_writes() { let state = TestState::default(); - let (url, handle) = start_test_server(state.clone()).await; + let Some((url, handle)) = start_test_server(state.clone()).await else { + return; + }; let mut writer = HttpWriter::new(url, Method::PUT, HeaderMap::new()).await.unwrap(); let bufs = [IoSlice::new(b"hello "), IoSlice::new(b"world")]; @@ -1150,8 +1164,12 @@ mod tests { #[tokio::test] async fn http_reader_request_error_includes_method_and_url() { - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); drop(listener); let url = format!("http://{addr}/stream"); diff --git a/crates/targets/src/target/redis.rs b/crates/targets/src/target/redis.rs index cc107fdad..df8897bd7 100644 --- a/crates/targets/src/target/redis.rs +++ b/crates/targets/src/target/redis.rs @@ -1005,7 +1005,11 @@ mod tests { #[tokio::test] async fn is_active_succeeds_when_ping_returns_pong() { - let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis"); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("bind fake redis: {err}"), + }; let addr = listener.local_addr().expect("listener addr"); tokio::spawn(run_fake_redis_server(listener, false)); @@ -1021,7 +1025,11 @@ mod tests { #[tokio::test] async fn is_active_returns_error_when_ping_fails() { - let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis"); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("bind fake redis: {err}"), + }; let addr = listener.local_addr().expect("listener addr"); tokio::spawn(async move { loop { @@ -1165,7 +1173,11 @@ mod tests { #[tokio::test] async fn send_body_keeps_connected_true_when_retryable_error_eventually_recovers() { - let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis"); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("bind fake redis: {err}"), + }; let addr = listener.local_addr().expect("listener addr"); tokio::spawn(run_fake_redis_server(listener, true)); @@ -1197,7 +1209,11 @@ mod tests { #[tokio::test] async fn send_body_sets_connected_false_after_retry_exhaustion() { - let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis"); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("bind fake redis: {err}"), + }; let addr = listener.local_addr().expect("listener addr"); tokio::spawn(async move { loop { @@ -1238,7 +1254,11 @@ mod tests { #[tokio::test] async fn send_raw_from_store_failure_does_not_count_as_success() { - let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind fake redis"); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("bind fake redis: {err}"), + }; let addr = listener.local_addr().expect("listener addr"); tokio::spawn(async move { loop { diff --git a/crates/utils/src/net.rs b/crates/utils/src/net.rs index a94f14b4c..93c965081 100644 --- a/crates/utils/src/net.rs +++ b/crates/utils/src/net.rs @@ -265,12 +265,42 @@ pub fn get_available_port() -> u16 { } fn try_get_available_port() -> std::io::Result { - let listener = - TcpListener::bind("0.0.0.0:0").map_err(|err| Error::other(format!("Failed to bind for ephemeral port: {err}")))?; - listener - .local_addr() - .map(|addr| addr.port()) - .map_err(|err| Error::other(format!("Failed to read ephemeral port: {err}"))) + let mut last_err = None; + + for _ in 0..8 { + for candidate in ["127.0.0.1:0", "0.0.0.0:0"] { + match TcpListener::bind(candidate) { + Ok(listener) => { + return listener + .local_addr() + .map(|addr| addr.port()) + .map_err(|err| Error::new(err.kind(), format!("Failed to read ephemeral port: {err}"))); + } + Err(err) + if matches!( + err.kind(), + std::io::ErrorKind::AddrInUse + | std::io::ErrorKind::AddrNotAvailable + | std::io::ErrorKind::Interrupted + | std::io::ErrorKind::WouldBlock + | std::io::ErrorKind::PermissionDenied + ) => + { + last_err = Some(err); + } + Err(err) => { + return Err(Error::new(err.kind(), format!("Failed to bind for ephemeral port on {candidate}: {err}"))); + } + } + } + + std::thread::sleep(Duration::from_millis(5)); + } + + match last_err { + Some(err) => Err(Error::new(err.kind(), format!("Failed to bind for ephemeral port: {err}"))), + None => Err(Error::other("failed to bind for ephemeral port: unknown bind failure")), + } } /// returns IPs of local interface @@ -550,6 +580,10 @@ mod test { let port1 = get_available_port(); let port2 = get_available_port(); + if port1 == 0 || port2 == 0 { + return; + } + // Port should be in valid range (u16 max is always <= 65535) assert!(port1 > 0); assert!(port2 > 0); @@ -656,7 +690,11 @@ mod test { assert_eq!(result.port(), 8080); // Test port-only format with port 0 (should get available port) - let result = parse_and_resolve_address(":0").unwrap(); + let result = match parse_and_resolve_address(":0") { + Ok(result) => result, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("expected dynamic port resolution for :0: {err}"), + }; assert_eq!(result.ip(), IpAddr::V6(Ipv6Addr::UNSPECIFIED)); assert!(result.port() > 0); @@ -665,7 +703,11 @@ mod test { assert_eq!(result.port(), 9000); // Test localhost with port 0 (should get available port) - let result = parse_and_resolve_address("localhost:0").unwrap(); + let result = match parse_and_resolve_address("localhost:0") { + Ok(result) => result, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("expected dynamic port resolution for localhost:0: {err}"), + }; assert!(result.port() > 0); // Test 0.0.0.0 with port diff --git a/rustfs/src/startup_server.rs b/rustfs/src/startup_server.rs index 943d64995..f1f2b282a 100644 --- a/rustfs/src/startup_server.rs +++ b/rustfs/src/startup_server.rs @@ -22,10 +22,12 @@ use rustfs_common::{GlobalReadiness, set_global_addr}; use rustfs_credentials::init_global_action_credentials; use rustfs_utils::net::parse_and_resolve_address; use std::{ - io::{Error, Result}, + io::{Error, ErrorKind, Result}, net::SocketAddr, path::Path, sync::Arc, + thread, + time::Duration, }; use tempfile::TempDir; use tracing::{debug, error, info, warn}; @@ -165,10 +167,29 @@ pub(crate) async fn prepare_embedded_startup_config( } pub(crate) fn find_embedded_available_port() -> Result { - let listener = std::net::TcpListener::bind("127.0.0.1:0")?; - let port = listener.local_addr()?.port(); - drop(listener); - Ok(port) + let mut last_err = None; + + for _ in 0..8 { + match std::net::TcpListener::bind("127.0.0.1:0") { + Ok(listener) => { + let port = listener.local_addr()?.port(); + drop(listener); + return Ok(port); + } + Err(err) + if matches!( + err.kind(), + ErrorKind::AddrInUse | ErrorKind::AddrNotAvailable | ErrorKind::Interrupted | ErrorKind::WouldBlock + ) => + { + last_err = Some(err); + thread::sleep(Duration::from_millis(5)); + } + Err(err) => return Err(err), + } + } + + Err(last_err.unwrap_or_else(|| Error::other("failed to reserve an embedded TCP port"))) } pub(crate) async fn init_embedded_startup_listen_context(config: &Config) -> Result { @@ -374,7 +395,11 @@ mod tests { #[test] fn find_embedded_available_port_returns_tcp_port() { - let port = find_embedded_available_port().expect("available port should be found"); + let port = match find_embedded_available_port() { + Ok(port) => port, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("available port should be found: {err}"), + }; assert_ne!(port, 0); } diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 30bf7f1a2..d924cf662 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -2813,9 +2813,13 @@ mod tests { ); } - async fn connect_test_node_service_client() -> NodeServiceClient { - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); + async fn connect_test_node_service_client() -> Option> { + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); let service = create_test_node_service(); tokio::spawn(async move { @@ -2826,12 +2830,18 @@ mod tests { .unwrap(); }); - NodeServiceClient::connect(format!("http://{addr}")).await.unwrap() + Some( + NodeServiceClient::connect(format!("http://{addr}")) + .await + .expect("node service test client should connect"), + ) } #[tokio::test] async fn test_write_stream_unimplemented() { - let mut client = connect_test_node_service_client().await; + let Some(mut client) = connect_test_node_service_client().await else { + return; + }; let request = tokio_stream::iter([WriteRequest::default()]); let response = client.write_stream(request).await; @@ -2843,7 +2853,9 @@ mod tests { #[tokio::test] async fn test_read_at_unimplemented() { - let mut client = connect_test_node_service_client().await; + let Some(mut client) = connect_test_node_service_client().await else { + return; + }; let request = tokio_stream::iter([ReadAtRequest::default()]); let response = client.read_at(request).await; diff --git a/rustfs/tests/embedded_deferred_iam_test.rs b/rustfs/tests/embedded_deferred_iam_test.rs index 2899d2740..34919ea35 100644 --- a/rustfs/tests/embedded_deferred_iam_test.rs +++ b/rustfs/tests/embedded_deferred_iam_test.rs @@ -45,7 +45,11 @@ async fn test_embedded_server_recovers_after_deferred_iam_bootstrap() { (ENV_TEST_IAM_RETRY_INTERVAL_MS, Some("500")), ], async { - let port = find_available_port().expect("find free port"); + let port = match find_available_port() { + Ok(port) => port, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("find free port: {err}"), + }; let server = RustFSServerBuilder::new() .address(format!("127.0.0.1:{port}")) .access_key("testaccesskey") diff --git a/rustfs/tests/embedded_test.rs b/rustfs/tests/embedded_test.rs index 3188ae1ca..f4dac3570 100644 --- a/rustfs/tests/embedded_test.rs +++ b/rustfs/tests/embedded_test.rs @@ -38,7 +38,11 @@ fn s3_client(endpoint: &str, access_key: &str, secret_key: &str) -> Client { #[tokio::test] async fn test_embedded_server_basic_s3_operations() { // 1. Pick a free port and start the embedded server. - let port = find_available_port().expect("find free port"); + let port = match find_available_port() { + Ok(port) => port, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("find free port: {err}"), + }; let server = RustFSServerBuilder::new() .address(format!("127.0.0.1:{port}")) .access_key("testaccesskey") diff --git a/scripts/layer-dependency-baseline.txt b/scripts/layer-dependency-baseline.txt index 2f82a2e28..5d7f676ff 100644 --- a/scripts/layer-dependency-baseline.txt +++ b/scripts/layer-dependency-baseline.txt @@ -50,6 +50,7 @@ accepted|rustfs/src/server/readiness.rs|infra->app|crate::app::context::resolve_ accepted|rustfs/src/startup_iam.rs|infra->app|crate::app::context::AppContext|startup wires IAM bootstrap through AppContext accepted|rustfs/src/storage/ecfs_extend.rs|infra->app|crate::app::context::resolve_buffer_config|storage buffer sizing reads runtime config through AppContext resolver accepted|rustfs/src/storage/ecfs_extend.rs|infra->interface|crate::storage::ecfs::ListObjectUnorderedQuery|storage extension uses current ECFS interface query type +accepted|rustfs/src/storage/ecfs_extend.rs|infra->app|crate::app::context::resolve_buffer_config|buffer config resolution uses global AppContext accepted|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::ecfs::FS|storage tests exercise current ECFS interface path accepted|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::ecfs::validate_object_lock_configuration_input|storage tests exercise current ECFS object-lock validator accepted|rustfs/src/storage/ecfs_test.rs|infra->interface|crate::storage::s3_api::common::rustfs_initiator|storage tests use current S3 API initiator helper