From 1bbfa71b11c83aabfd9d4c9795142de619c4dc07 Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 2 Sep 2026 15:33:55 +0800 Subject: [PATCH] fix(ecstore): preserve buckets after pool expansion (#7040) * fix(ecstore): preserve buckets after pool expansion * fix(ecstore): scope bucket operations by erasure set * fix(ecstore): preserve bucket metadata load errors --- crates/ecstore/src/bucket/metadata_sys.rs | 81 ++- .../ecstore/src/cluster/rpc/peer_s3_client.rs | 2 +- crates/ecstore/src/store/bucket.rs | 593 +++++++++++++++--- 3 files changed, 559 insertions(+), 117 deletions(-) diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index d9f1892c1..a05114a01 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -21,7 +21,7 @@ use crate::bucket::bucket_target_sys::BucketTargetSys; use crate::bucket::metadata::{load_bucket_metadata_parse, load_bucket_metadata_parse_with_presence}; use crate::bucket::utils::is_meta_bucketname; use crate::disk::RUSTFS_META_BUCKET; -use crate::error::{Error, Result, is_err_bucket_not_found}; +use crate::error::{Error, Result, is_err_bucket_not_found, is_err_strict_volume_not_found}; use crate::runtime::sources as runtime_sources; use crate::storage_api_contracts::heal::HealOperations as _; use crate::storage_api_contracts::namespace::NamespaceLocking as _; @@ -659,13 +659,12 @@ async fn acquire_config_write_guard_for_incarnation( async { match metadata_sys .api - .peer_sys - .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) .await { Ok(_) => Ok(()), - Err(crate::disk::error::Error::VolumeNotFound) => Err(Error::BucketNotFound(bucket.to_string())), - Err(err) => Err(err.into()), + Err(err) if is_err_strict_volume_not_found(&err) => Err(Error::BucketNotFound(bucket.to_string())), + Err(err) => Err(err), } }, ), @@ -1274,7 +1273,6 @@ pub struct BucketMetadataSys { /// name floods while avoiding repeated namespace and erasure reads. missing_buckets: moka::future::Cache, api: Arc, - initialized: Arc>, } impl BucketMetadataSys { @@ -1303,7 +1301,6 @@ impl BucketMetadataSys { .time_to_live(MISSING_BUCKET_TTL) .build(), api, - initialized: Arc::new(RwLock::new(false)), } } @@ -1364,13 +1361,12 @@ impl BucketMetadataSys { await_bucket_namespace_operation(Some(namespace_guard), bucket, operation, async { match self .api - .peer_sys - .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) .await { Ok(_) => Ok(true), - Err(crate::disk::error::Error::VolumeNotFound) => Ok(false), - Err(err) => Err(err.into()), + Err(Error::VolumeNotFound) => Ok(false), + Err(err) => Err(err), } }) .await @@ -1380,9 +1376,15 @@ impl BucketMetadataSys { let _ = self.init_internal(buckets).await; } async fn init_internal(&self, buckets: Vec) -> Result<()> { - let count = runtime_sources::endpoint_erasure_set_count() - .map(|count| count * 10) - .ok_or_else(|| Error::other("endpoint pools not initialized"))?; + let count = self + .api + .pools + .iter() + .map(|pool| pool.disk_set.len()) + .sum::() + .checked_mul(10) + .filter(|count| *count != 0) + .ok_or_else(|| Error::other("bucket metadata store has no erasure sets"))?; let mut failed_buckets: HashSet = HashSet::new(); let mut buckets = buckets.as_slice(); @@ -1400,9 +1402,6 @@ impl BucketMetadataSys { buckets = &buckets[count..] } - let mut initialized = self.initialized.write().await; - *initialized = true; - Ok(()) } @@ -1479,14 +1478,6 @@ impl BucketMetadataSys { expected: Option<&Arc>, namespace_guard: &rustfs_lock::NamespaceLockGuard, ) -> Result<()> { - await_bucket_namespace_operation( - Some(namespace_guard), - bucket, - "bucket metadata heal", - self.api.heal_bucket(bucket, &HealOpts::default()), - ) - .await?; - if !self .bucket_exists(bucket, namespace_guard, "bucket metadata existence check") .await? @@ -1506,6 +1497,20 @@ impl BucketMetadataSys { return Ok(()); } + await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "bucket metadata heal", + self.api.heal_bucket( + bucket, + &HealOpts { + recreate: true, + ..Default::default() + }, + ), + ) + .await?; + let (bm, persisted) = await_bucket_namespace_operation( Some(namespace_guard), bucket, @@ -1887,23 +1892,13 @@ impl BucketMetadataSys { "lazy metadata IO must start while the bucket namespace read lock is held" ); } - let (bm, persisted) = match await_bucket_namespace_operation( + let (bm, persisted) = await_bucket_namespace_operation( Some(&guard), bucket, "lazy bucket metadata load", Box::pin(load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true)), ) - .await - { - Ok(res) => res, - Err(err) => { - return if *self.initialized.read().await { - Err(Error::other("errBucketMetadataNotInitialized")) - } else { - Err(err) - }; - } - }; + .await?; let bm = Arc::new(bm); @@ -1914,11 +1909,9 @@ impl BucketMetadataSys { "lazy bucket metadata existence check", Box::pin(async { self.api - .peer_sys - .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) .await .map(|_| ()) - .map_err(Into::into) }), ) .await?; @@ -2197,10 +2190,8 @@ impl BucketMetadataSys { "legacy bucket metadata existence check", async { self.api - .peer_sys - .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) .await - .map_err(crate::error::StorageError::from) }, ) .await @@ -2302,10 +2293,8 @@ impl BucketMetadataSys { "bucket metadata snapshot existence check", async { self.api - .peer_sys - .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) .await - .map_err(crate::error::StorageError::from) }, ) .await diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index 8e5ebcecc..6a929113e 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -190,7 +190,7 @@ fn install_heal_bucket_pre_mutation_barrier() -> Arc kind.as_str()).increment(1); @@ -49,8 +50,8 @@ fn record_bucket_delete_blocker(bucket: &str, kind: BucketDeleteBlockerKind, res ); } -fn scanner_bucket_list_set_concurrency(set_count: usize) -> usize { - set_count.clamp(1, SCANNER_BUCKET_LIST_SET_CONCURRENCY) +fn bucket_list_set_concurrency(set_count: usize, max_concurrency: usize) -> usize { + set_count.clamp(1, max_concurrency.max(1)) } fn should_override_created_from_metadata(created: OffsetDateTime) -> bool { @@ -223,6 +224,14 @@ async fn bucket_delete_local_blocker( } impl ECStore { + fn bucket_sets(&self) -> impl Iterator)> + '_ { + self.pools.iter().flat_map(|pool| { + pool.disk_set + .iter() + .map(|set| (set.pool_index, set.set_index, Arc::clone(set))) + }) + } + pub async fn get_bucket_metadata(&self, bucket: &str) -> Result> { let sys = metadata_sys::require_bucket_metadata_sys_in(&self.ctx)?; sys.read().await.get(bucket).await @@ -341,20 +350,96 @@ impl ECStore { async fn mark_bucket_deleted(&self, bucket: &str) -> Result<()> { let marker_volume = bucket_deleted_marker_volume(bucket); - self.peer_sys - .make_bucket( - marker_volume.as_str(), - &MakeBucketOptions { - force_create: true, - ..Default::default() - }, - ) - .await - .map_err(|e| to_object_err(e.into(), vec![bucket]))?; + self.make_bucket_on_sets( + marker_volume.as_str(), + &MakeBucketOptions { + force_create: true, + ..Default::default() + }, + ) + .await + .map_err(|err| to_object_err(err, vec![bucket]))?; Ok(()) } + async fn make_bucket_on_sets(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> { + let results = futures::future::join_all( + self.bucket_sets() + .map(|(_, _, set)| async move { set.make_bucket(bucket, opts).await }), + ) + .await; + if results.is_empty() { + return Err(StorageError::ErasureWriteQuorum); + } + + let mut bucket_exists_error = None; + let mut first_hard_error = None; + for result in results { + let Err(err) = result else { + continue; + }; + if matches!(err, StorageError::VolumeExists) || is_err_bucket_exists(&err) { + if bucket_exists_error.is_none() { + bucket_exists_error = Some(err); + } + } else if first_hard_error.is_none() { + first_hard_error = Some(err); + } + } + + match first_hard_error.or(bucket_exists_error) { + Some(err) => Err(err), + None => Ok(()), + } + } + + async fn delete_bucket_on_sets(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { + let results = futures::future::join_all( + self.bucket_sets() + .map(|(_, _, set)| async move { set.delete_bucket(bucket, opts).await }), + ) + .await; + if results.is_empty() { + return Err(StorageError::ErasureWriteQuorum); + } + + let mut deleted = false; + let mut first_error = None; + for result in results { + match result { + Ok(()) => deleted = true, + Err(err) if is_err_strict_volume_not_found(&err) => {} + Err(err) if first_error.is_none() => first_error = Some(err), + Err(_) => {} + } + } + + if let Some(delete_error) = first_error { + if !opts.no_recreate { + let rollback_opts = MakeBucketOptions { + force_create: true, + no_lock: true, + ..Default::default() + }; + if let Err(rollback_error) = self.make_bucket_on_sets(bucket, &rollback_opts).await { + warn!( + event = EVENT_BUCKET_DELETE_ROLLBACK_FAILED, + component = "ecstore", + subsystem = "bucket", + bucket, + deletion_error = ?delete_error, + rollback_error = ?rollback_error, + "Bucket deletion rollback could not restore every bucket volume" + ); + } + } + return Err(delete_error); + } + + if deleted { Ok(()) } else { Err(StorageError::VolumeNotFound) } + } + async fn cleanup_deleted_bucket_metadata( &self, bucket: &str, @@ -424,10 +509,9 @@ impl ECStore { }; if let Err(err) = await_bucket_lifecycle_operation(lifecycle_guard, namespace_guard, bucket, "failed bucket creation rollback", async { - self.peer_sys - .delete_bucket(bucket, &rollback_opts) + self.delete_bucket_on_sets(bucket, &rollback_opts) .await - .map_err(|rollback_err| to_object_err(rollback_err.into(), vec![bucket])) + .map_err(|rollback_err| to_object_err(rollback_err, vec![bucket])) }) .await { @@ -494,10 +578,9 @@ impl ECStore { None }; - let existing_bucket_info = match self.peer_sys.get_bucket_info(bucket, &BucketOptions::default()).await { + let existing_bucket_info = match self.get_bucket_info_from_sets(bucket, &BucketOptions::default()).await { Ok(info) => Some(info), Err(err) => { - let err: StorageError = err.into(); if is_err_bucket_not_found(&err) { None } else { @@ -586,10 +669,9 @@ impl ECStore { bucket, "bucket creation metadata transaction", async { - self.peer_sys - .make_bucket(bucket, opts) + self.make_bucket_on_sets(bucket, opts) .await - .map_err(|err| to_object_err(err.into(), vec![bucket])) + .map_err(|err| to_object_err(err, vec![bucket])) }, ), ) @@ -608,7 +690,7 @@ impl ECStore { { warn!("best-effort bucket heal after BucketExists failed: {heal_err}"); } - if !is_err_bucket_exists(&err) && ns_guard.as_ref().is_none_or(|guard| !guard.is_lock_lost()) { + if confirmed_missing && !is_err_bucket_exists(&err) && ns_guard.as_ref().is_none_or(|guard| !guard.is_lock_lost()) { error!("make bucket failed: {err}"); self.rollback_failed_bucket_creation(bucket, bucket_lifecycle_guard.as_ref(), ns_guard.as_ref()) .await; @@ -653,9 +735,43 @@ impl ECStore { Ok(()) } + #[instrument(skip(self))] + pub(crate) async fn get_bucket_info_from_sets(&self, bucket: &str, opts: &BucketOptions) -> Result { + // One host may participate in several pools after expansion. Resolve the + // namespace against each erasure set so disks from different pools can + // never be combined into one bucket quorum. + // Bucket validation is request-path IO. Keep the previous peer fanout's + // latency shape by probing every set concurrently; scanner listings use + // a separate bounded path below because they run continuously. + let mut scoped_results = + futures::future::join_all(self.bucket_sets().map(|(pool_index, set_index, set)| async move { + (pool_index, set_index, set.get_bucket_info(bucket, opts).await) + })) + .await; + scoped_results.sort_unstable_by_key(|(pool_index, set_index, _)| (*pool_index, *set_index)); + + let mut first_info = None; + let mut first_error = None; + for (_, _, result) in scoped_results { + match result { + Ok(info) if first_info.is_none() => first_info = Some(info), + Ok(_) => {} + Err(err) if is_err_strict_volume_not_found(&err) => {} + Err(err) if first_error.is_none() => first_error = Some(err), + Err(_) => {} + } + } + + if let Some(err) = first_error { + return Err(err); + } + + first_info.ok_or(Error::VolumeNotFound) + } + #[instrument(skip(self))] pub(super) async fn handle_get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result { - let mut info = self.peer_sys.get_bucket_info(bucket, opts).await?; + let mut info = self.get_bucket_info_from_sets(bucket, opts).await?; if let Ok(sys) = metadata_sys::get_in(&self.ctx, bucket).await { if should_override_created_from_metadata(sys.created) { @@ -671,36 +787,24 @@ impl ECStore { #[instrument(skip(self))] pub(super) async fn handle_list_bucket(&self, opts: &BucketOptions) -> Result> { // TODO(backlog): support cached bucket listing via opts.cached - - let mut buckets = self.peer_sys.list_bucket(opts).await?; - - if !opts.no_metadata { - for bucket in buckets.iter_mut() { - if let Ok(created) = metadata_sys::created_at_in(&self.ctx, &bucket.name).await - && should_override_created_from_metadata(created) - { - bucket.created = Some(created); - } - } - } - Ok(buckets) + Ok(self.list_bucket_from_sets(opts, usize::MAX).await?.buckets) } pub async fn list_bucket_for_scanner(&self, opts: &BucketOptions) -> Result { - let sets = self - .pools - .iter() - .flat_map(|pool| { - pool.disk_set - .iter() - .map(|set| (set.pool_index, set.set_index, Arc::clone(set))) - }) - .collect::>(); - let set_count = sets.len(); + self.list_bucket_from_sets(opts, SCANNER_BUCKET_LIST_SET_CONCURRENCY).await + } + + async fn list_bucket_from_sets( + &self, + opts: &BucketOptions, + max_concurrency: usize, + ) -> Result { + let set_count = self.pools.iter().map(|pool| pool.disk_set.len()).sum(); + let concurrency = bucket_list_set_concurrency(set_count, max_concurrency); let deleted = opts.deleted; let cached = opts.cached; let no_metadata = opts.no_metadata; - let mut set_listings = stream::iter(sets.into_iter().map(move |(pool_index, set_index, set)| { + let mut set_listings = stream::iter(self.bucket_sets().map(move |(pool_index, set_index, set)| { let opts = BucketOptions { deleted, cached, @@ -712,7 +816,7 @@ impl ECStore { .map(|(buckets, complete)| (pool_index, set_index, buckets, complete)) } })) - .buffer_unordered(scanner_bucket_list_set_concurrency(set_count)); + .buffer_unordered(concurrency); let mut topology_complete = set_count != 0; let mut bucket_map = BTreeMap::::new(); let mut scoped_buckets = Vec::with_capacity(set_count); @@ -817,16 +921,15 @@ impl ECStore { let sr_purge = opts.srdelete_op == SRBucketDeleteOp::Purge; let sr_delete = sr_mark_delete || sr_purge; let mut delete_opts = opts.clone(); - let bucket_exists = match self.peer_sys.get_bucket_info(bucket, &BucketOptions::default()).await { + let bucket_exists = match self.get_bucket_info_from_sets(bucket, &BucketOptions::default()).await { Ok(_) => true, Err(err) => { - let storage_err: StorageError = err.into(); - if is_err_strict_volume_not_found(&storage_err) && sr_delete { + if is_err_strict_volume_not_found(&err) && sr_delete { false - } else if is_err_strict_volume_not_found(&storage_err) { + } else if is_err_strict_volume_not_found(&err) { return Err(StorageError::BucketNotFound(bucket.to_string())); } else { - return Err(to_object_err(storage_err, vec![bucket])); + return Err(to_object_err(err, vec![bucket])); } } }; @@ -845,6 +948,9 @@ impl ECStore { delete_opts.force_if_empty = true; } + #[cfg(test)] + crate::cluster::rpc::peer_s3_client::pause_after_delete_bucket_empty_scan().await; + if sr_mark_delete { await_bucket_lifecycle_operation( bucket_lifecycle_guard.as_ref(), @@ -861,10 +967,9 @@ impl ECStore { bucket, "physical bucket deletion", run_physical_bucket_deletion(ns_guard.as_ref(), bucket, async { - self.peer_sys - .delete_bucket(bucket, &delete_opts) + self.delete_bucket_on_sets(bucket, &delete_opts) .await - .map_err(|err| to_object_err(err.into(), vec![bucket])) + .map_err(|err| to_object_err(err, vec![bucket])) }), ) .await; @@ -905,14 +1010,14 @@ mod tests { BUCKET_DELETE_DIAGNOSTIC_MAX_ELAPSED, BUCKET_DELETE_DIAGNOSTIC_MAX_ENTRIES, BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES, BucketDeleteBlockerKind, BucketDeleteDiagnosticBudget, SCANNER_BUCKET_LIST_SET_CONCURRENCY, await_bucket_namespace_operation, bucket_delete_metadata_cleanup_prefixes, bucket_deleted_marker_prefix, - bucket_deleted_marker_volume, run_bucket_usage_cleanup, run_physical_bucket_deletion, scan_metadata_less_residue, - scan_metadata_less_residue_with_budget, scanner_bucket_list_set_concurrency, should_override_created_from_metadata, + bucket_deleted_marker_volume, bucket_list_set_concurrency, run_bucket_usage_cleanup, run_physical_bucket_deletion, + scan_metadata_less_residue, scan_metadata_less_residue_with_budget, should_override_created_from_metadata, validate_table_bucket_delete_allowed, }; use crate::bucket::metadata::table_bucket_catalog_metadata_prefix; use crate::bucket::metadata_sys; use crate::cluster::rpc::peer_s3_client::install_delete_bucket_empty_scan_barrier; - use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE}; + use crate::disk::{BUCKET_META_PREFIX, DiskAPI, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE}; use crate::error::StorageError; use crate::object_api::{ObjectOptions, PutObjReader}; use crate::runtime::instance::InstanceContext; @@ -1217,8 +1322,8 @@ mod tests { .clone() } - async fn setup_multi_pool_scanner_listing_test_env() -> (tempfile::TempDir, Arc) { - let temp_dir = tempfile::tempdir().expect("multi-pool scanner test directory should be created"); + async fn setup_multi_pool_bucket_test_env() -> (tempfile::TempDir, Arc) { + let temp_dir = tempfile::tempdir().expect("multi-pool bucket test directory should be created"); let mut pools = Vec::new(); for pool_index in 0..2 { let mut endpoints = Vec::new(); @@ -1226,7 +1331,7 @@ mod tests { let disk_path = temp_dir.path().join(format!("pool{pool_index}-disk{disk_index}")); tokio::fs::create_dir_all(&disk_path) .await - .expect("multi-pool scanner test disk should be created"); + .expect("multi-pool bucket test disk should be created"); let mut endpoint = Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse"); endpoint.set_pool_index(pool_index); @@ -1239,13 +1344,14 @@ mod tests { set_count: 1, drives_per_set: 4, endpoints: Endpoints::from(endpoints), - cmd_line: format!("scanner-listing-pool-{pool_index}"), + cmd_line: format!("bucket-test-pool-{pool_index}"), platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH), }); } let endpoint_pools = EndpointServerPools(pools); let instance_ctx = Arc::new(InstanceContext::new()); + instance_ctx.set_endpoints(endpoint_pools.clone()); init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone()) .await .expect("multi-pool local disks should initialize"); @@ -1269,9 +1375,47 @@ mod tests { (temp_dir, ecstore) } + async fn take_set_disks_offline( + ecstore: &ECStore, + set: &Arc, + disk_indexes: &[usize], + ) -> Vec<(usize, crate::disk::DiskStore)> { + let offline = { + let mut disks = set.disks.write().await; + disk_indexes + .iter() + .map(|index| (*index, disks[*index].take().expect("fault-injection disk should start online"))) + .collect::>() + }; + let local_disk_map = ecstore.ctx.local_disk_map(); + let mut local_disks = local_disk_map.write().await; + for (_, disk) in &offline { + local_disks.insert(disk.endpoint().to_string(), None); + } + offline + } + + async fn restore_set_disks( + ecstore: &ECStore, + set: &Arc, + offline: Vec<(usize, crate::disk::DiskStore)>, + ) { + { + let local_disk_map = ecstore.ctx.local_disk_map(); + let mut local_disks = local_disk_map.write().await; + for (_, disk) in &offline { + local_disks.insert(disk.endpoint().to_string(), Some(Arc::clone(disk))); + } + } + let mut disks = set.disks.write().await; + for (index, disk) in offline { + assert!(disks[index].replace(disk).is_none(), "fault-injection disk slot should remain empty"); + } + } + #[tokio::test] async fn request_metadata_methods_fail_closed_before_instance_initialization() { - let (_temp_dir, store) = setup_multi_pool_scanner_listing_test_env().await; + let (_temp_dir, store) = setup_multi_pool_bucket_test_env().await; let expected = "bucket metadata sys not initialized for this instance"; let errors = [ @@ -1405,10 +1549,14 @@ mod tests { } #[test] - fn scanner_bucket_listing_bounds_set_fanout() { - assert_eq!(scanner_bucket_list_set_concurrency(0), 1); - assert_eq!(scanner_bucket_list_set_concurrency(2), 2); - assert_eq!(scanner_bucket_list_set_concurrency(100), SCANNER_BUCKET_LIST_SET_CONCURRENCY); + fn bucket_listing_selects_request_and_scanner_fanout() { + assert_eq!(bucket_list_set_concurrency(0, SCANNER_BUCKET_LIST_SET_CONCURRENCY), 1); + assert_eq!(bucket_list_set_concurrency(2, SCANNER_BUCKET_LIST_SET_CONCURRENCY), 2); + assert_eq!( + bucket_list_set_concurrency(100, SCANNER_BUCKET_LIST_SET_CONCURRENCY), + SCANNER_BUCKET_LIST_SET_CONCURRENCY + ); + assert_eq!(bucket_list_set_concurrency(100, usize::MAX), 100); } #[tokio::test] @@ -1613,10 +1761,315 @@ mod tests { assert!(oversized_scan.diagnostic_bytes_read <= BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES); } + #[tokio::test] + #[serial] + async fn bucket_namespace_reads_keep_pre_expansion_bucket_visible() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("pre-expansion-{}", Uuid::new_v4().simple()); + let object = "existing-object"; + ecstore.pools[0].disk_set[0] + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created in the original pool only"); + let mut reader = PutObjReader::from_vec(b"object written before pool expansion".to_vec()); + ecstore.pools[0] + .put_object(&bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("object should be written in the original pool only"); + + let buckets = ecstore + .list_bucket(&BucketOptions::default()) + .await + .expect("S3 bucket listing should remain available after adding an empty pool"); + assert!(buckets.iter().any(|entry| entry.name == bucket)); + + let info = ecstore + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("bucket validation should accept a bucket present in the original pool"); + assert_eq!(info.name, bucket); + + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), vec![bucket.clone()]).await; + let listed = ecstore + .clone() + .list_objects_v2(&bucket, "", None, None, 1000, false, None, false) + .await + .expect("ListObjectsV2 should remain available after adding an empty pool"); + assert!(listed.objects.iter().any(|entry| entry.name == object)); + ecstore.pools[1].disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("metadata initialization should heal the bucket volume into the expansion pool"); + + let object_info = ecstore + .get_object_info(&bucket, object, &ObjectOptions::default()) + .await + .expect("an object in the original pool should remain readable after expansion"); + assert_eq!(object_info.name, object); + } + + #[tokio::test] + #[serial] + async fn bucket_creation_accepts_exact_quorum_in_every_set() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("exact-set-quorum-{}", Uuid::new_v4().simple()); + let original_set = Arc::clone(&ecstore.pools[0].disk_set[0]); + let offline = take_set_disks_offline(&ecstore, &original_set, &[0]).await; + + ecstore + .make_bucket_on_sets(&bucket, &MakeBucketOptions::default()) + .await + .expect("three of four disks should satisfy the original set quorum"); + + restore_set_disks(&ecstore, &original_set, offline).await; + for pool in &ecstore.pools { + pool.disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("every set should expose a bucket created at its exact quorum"); + } + } + + #[tokio::test] + #[serial] + async fn bucket_creation_rejects_cross_pool_quorum_subsidy() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), Vec::new()).await; + let bucket = format!("create-set-quorum-minus-one-{}", Uuid::new_v4().simple()); + let original_set = Arc::clone(&ecstore.pools[0].disk_set[0]); + let offline = take_set_disks_offline(&ecstore, &original_set, &[0, 1]).await; + + let err = ecstore + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect_err("a healthy expansion pool must not subsidize a failed original set"); + restore_set_disks(&ecstore, &original_set, offline).await; + assert!( + matches!(&err, StorageError::InsufficientWriteQuorum(name, object) if name == &bucket && object.is_empty()), + "unexpected error: {err}" + ); + for pool in &ecstore.pools { + assert_eq!( + pool.disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect_err("failed creation must not leave an authoritative bucket"), + StorageError::VolumeNotFound + ); + } + } + + #[tokio::test] + #[serial] + async fn bucket_creation_prioritizes_hard_set_error_over_existing_bucket() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), Vec::new()).await; + let bucket = format!("create-hard-error-{}", Uuid::new_v4().simple()); + ecstore.pools[0].disk_set[0] + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("the original pool should contain the bucket"); + let expansion_set = Arc::clone(&ecstore.pools[1].disk_set[0]); + let offline = take_set_disks_offline(&ecstore, &expansion_set, &[0, 1]).await; + + let err = ecstore + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect_err("a hard set failure must outrank BucketExists from another set"); + restore_set_disks(&ecstore, &expansion_set, offline).await; + assert!( + matches!(&err, StorageError::InsufficientWriteQuorum(name, object) if name == &bucket && object.is_empty()), + "unexpected error: {err}" + ); + ecstore.pools[0].disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("failure in another set must not roll back a pre-existing bucket"); + } + + #[tokio::test] + #[serial] + async fn bucket_delete_rejects_cross_pool_quorum_subsidy_before_mutation() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), Vec::new()).await; + let bucket = format!("delete-set-quorum-minus-one-{}", Uuid::new_v4().simple()); + ecstore + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created in every set before deletion"); + let original_set = Arc::clone(&ecstore.pools[0].disk_set[0]); + let offline = take_set_disks_offline(&ecstore, &original_set, &[0, 1, 2]).await; + + let err = ecstore + .delete_bucket( + &bucket, + &DeleteBucketOptions { + force: true, + ..Default::default() + }, + ) + .await + .expect_err("a healthy expansion pool must not subsidize a failed original set deletion"); + restore_set_disks(&ecstore, &original_set, offline).await; + assert!( + matches!(&err, StorageError::InsufficientWriteQuorum(name, object) if name == &bucket && object.is_empty()), + "unexpected error: {err}" + ); + for pool in &ecstore.pools { + pool.disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("failed deletion preflight must preserve every bucket volume"); + } + let (_, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&ecstore.ctx, &bucket) + .await + .expect("failed deletion must preserve bucket metadata"); + assert!(persisted); + } + + #[tokio::test] + #[serial] + async fn bucket_delete_rolls_back_sets_after_partial_failure() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("delete-set-rollback-{}", Uuid::new_v4().simple()); + ecstore + .make_bucket_on_sets(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created in every set before deletion"); + let original_set = Arc::clone(&ecstore.pools[0].disk_set[0]); + let offline = take_set_disks_offline(&ecstore, &original_set, &[0, 1]).await; + + let err = ecstore + .delete_bucket_on_sets(&bucket, &DeleteBucketOptions::default()) + .await + .expect_err("a failed set must fail the complete bucket deletion"); + restore_set_disks(&ecstore, &original_set, offline).await; + assert_eq!(err, StorageError::ErasureWriteQuorum); + for pool in &ecstore.pools { + pool.disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("delete rollback should restore the bucket namespace in every set"); + } + } + + #[tokio::test] + #[serial] + async fn bucket_config_update_accepts_pre_expansion_bucket() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), Vec::new()).await; + let bucket = format!("config-pre-expansion-{}", Uuid::new_v4().simple()); + ecstore + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created with authoritative metadata"); + ecstore.pools[1].disk_set[0] + .delete_bucket(&bucket, &DeleteBucketOptions::default()) + .await + .expect("the expansion pool bucket volume should be removed for the regression state"); + + let policy = br#"{"Version":"2012-10-17","Statement":[]}"#.to_vec(); + ecstore + .update_bucket_metadata_config(&bucket, crate::bucket::metadata::BUCKET_POLICY_CONFIG, policy.clone()) + .await + .expect("config mutation should validate existence per erasure set"); + let (stored, _) = ecstore + .get_bucket_policy_raw(&bucket) + .await + .expect("updated bucket policy should remain readable"); + assert_eq!(stored.as_bytes(), policy); + } + + #[tokio::test] + #[serial] + async fn bucket_metadata_init_recreates_pre_expansion_bucket_in_new_pool() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("metadata-expansion-{}", Uuid::new_v4().simple()); + ecstore.pools[0].disk_set[0] + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created in the original pool only"); + assert_eq!( + ecstore.pools[1].disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect_err("new pool should initially have no bucket volume"), + StorageError::VolumeNotFound + ); + + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), vec![bucket.clone()]).await; + + ecstore.pools[1].disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect("metadata initialization should recreate the bucket volume in the new pool"); + } + + #[tokio::test] + #[serial] + async fn bucket_metadata_init_does_not_recreate_stale_bucket_name() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("stale-expansion-{}", Uuid::new_v4().simple()); + + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), vec![bucket.clone()]).await; + + for pool in &ecstore.pools { + let err = pool.disk_set[0] + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect_err("a stale startup listing must not recreate a deleted bucket"); + assert_eq!(err, StorageError::VolumeNotFound); + } + } + + #[tokio::test] + #[serial] + async fn bucket_namespace_reads_report_missing_when_every_set_is_absent() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("absent-{}", Uuid::new_v4().simple()); + + let buckets = ecstore + .list_bucket(&BucketOptions::default()) + .await + .expect("an empty healthy namespace should remain listable"); + assert!(buckets.iter().all(|entry| entry.name != bucket)); + + let err = ecstore + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect_err("a bucket absent from every set must remain missing"); + assert_eq!(err, StorageError::VolumeNotFound); + } + + #[tokio::test] + #[serial] + async fn bucket_namespace_reads_fail_closed_when_any_set_loses_quorum() { + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("degraded-expansion-{}", Uuid::new_v4().simple()); + ecstore.pools[0].disk_set[0] + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created in the original pool only"); + ecstore.pools[1].disk_set[0].disks.write().await[0] = None; + ecstore.pools[1].disk_set[0].disks.write().await[1] = None; + + let list_err = ecstore + .list_bucket(&BucketOptions::default()) + .await + .expect_err("listing must not hide an unavailable expansion pool"); + assert_eq!(list_err, StorageError::ErasureWriteQuorum); + + let info_err = ecstore + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .expect_err("bucket validation must fail when an expansion pool is unavailable"); + assert_eq!(info_err, StorageError::ErasureWriteQuorum); + } + #[tokio::test] #[serial] async fn scanner_bucket_listing_unions_every_erasure_set() { - let (_temp_dir, ecstore) = setup_multi_pool_scanner_listing_test_env().await; + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; let bucket = format!("second-pool-only-{}", Uuid::new_v4().simple()); ecstore.pools[1].disk_set[0] .make_bucket(&bucket, &MakeBucketOptions::default()) @@ -1652,7 +2105,7 @@ mod tests { #[tokio::test] async fn scanner_bucket_listing_marks_degraded_set_incomplete() { - let (_temp_dir, ecstore) = setup_multi_pool_scanner_listing_test_env().await; + let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; let bucket = format!("degraded-set-{}", Uuid::new_v4().simple()); let set = &ecstore.pools[0].disk_set[0]; set.make_bucket(&bucket, &MakeBucketOptions::default()) @@ -1674,7 +2127,7 @@ mod tests { #[tokio::test] async fn scanner_bucket_listing_marks_divergent_disk_views_incomplete() { - let (temp_dir, ecstore) = setup_multi_pool_scanner_listing_test_env().await; + let (temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; let bucket = format!("divergent-set-{}", Uuid::new_v4().simple()); ecstore.pools[0].disk_set[0] .make_bucket(&bucket, &MakeBucketOptions::default())