From 1cb1b02b08b840e9235837f2c71ea30e7438a4db Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 29 Jul 2026 14:49:18 +0800 Subject: [PATCH] fix(ecstore): prevent deleted bucket recreation (#5380) * fix(ecstore): prevent deleted bucket recreation * fix(ecstore): reject empty metadata bucket names * fix(ecstore): avoid recursive bucket metadata lookup * fix(ecstore): bound lazy metadata future stack use * test(ecstore): create bucket before metadata reload --------- Co-authored-by: houseme --- crates/ecstore/src/bucket/metadata_sys.rs | 281 +++++++++++++++--- .../ecstore/src/cluster/rpc/peer_s3_client.rs | 45 ++- crates/ecstore/src/store/bucket.rs | 2 +- crates/ecstore/src/store/mod.rs | 1 + 4 files changed, 279 insertions(+), 50 deletions(-) diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index dbd9bb573..2a30a17be 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -23,7 +23,7 @@ use crate::error::{Error, Result, is_err_bucket_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 _; -use crate::store::ECStore; +use crate::store::{ECStore, await_bucket_namespace_operation}; use futures::future::join_all; use rustfs_common::heal_channel::HealOpts; use rustfs_policy::policy::BucketPolicy; @@ -37,13 +37,19 @@ use std::collections::HashSet; use std::time::Duration; use std::{collections::HashMap, sync::Arc}; use time::OffsetDateTime; -use tokio::sync::RwLock; +use tokio::sync::{Mutex, RwLock}; use tokio::time::sleep; use tokio_util::sync::CancellationToken; use tracing::{error, warn}; const BUCKET_METADATA_REFRESH_INTERVAL: Duration = Duration::from_secs(15 * 60); +#[derive(Clone, Copy)] +enum MetadataLoadMode { + Initial, + Refresh, +} + pub async fn init_bucket_metadata_sys(api: Arc, buckets: Vec) { // The metadata system is inherently per-store (it holds the store handle // and that store's bucket cache), so it lives on the store's own instance @@ -140,7 +146,8 @@ async fn refresh_buckets_metadata_once(sys: Arc>) { for chunk in buckets.chunks(count) { let sys = sys.read().await; - sys.concurrent_load(chunk, &mut failed_buckets).await; + sys.concurrent_load(chunk, &mut failed_buckets, MetadataLoadMode::Refresh) + .await; } if !failed_buckets.is_empty() { @@ -461,6 +468,9 @@ const ABSENT_BUCKET_METADATA_MAX_ENTRIES: u64 = 10_000; #[derive(Debug)] pub struct BucketMetadataSys { metadata_map: RwLock>>, + metadata_publish_lock: Mutex<()>, + #[cfg(test)] + lazy_load_lock_probe: std::sync::atomic::AtomicBool, /// Buckets recently observed to have no persisted metadata. Serving the /// fabricated default from here (instead of re-reading disk) keeps the /// per-request cost of repeated lookups for such names bounded — without @@ -476,6 +486,9 @@ impl BucketMetadataSys { pub fn new(api: Arc) -> Self { Self { metadata_map: RwLock::new(HashMap::new()), + metadata_publish_lock: Mutex::new(()), + #[cfg(test)] + lazy_load_lock_probe: std::sync::atomic::AtomicBool::new(false), absent_metadata: moka::future::Cache::builder() .max_capacity(ABSENT_BUCKET_METADATA_MAX_ENTRIES) .time_to_live(ABSENT_BUCKET_METADATA_TTL) @@ -502,11 +515,13 @@ impl BucketMetadataSys { loop { if buckets.len() < count { - self.concurrent_load(buckets, &mut failed_buckets).await; + self.concurrent_load(buckets, &mut failed_buckets, MetadataLoadMode::Initial) + .await; break; } - self.concurrent_load(&buckets[..count], &mut failed_buckets).await; + self.concurrent_load(&buckets[..count], &mut failed_buckets, MetadataLoadMode::Initial) + .await; buckets = &buckets[count..] } @@ -517,7 +532,7 @@ impl BucketMetadataSys { Ok(()) } - async fn concurrent_load(&self, buckets: &[String], failed_buckets: &mut HashSet) { + async fn concurrent_load(&self, buckets: &[String], failed_buckets: &mut HashSet, mode: MetadataLoadMode) { let mut futures = Vec::new(); for bucket in buckets.iter() { @@ -525,16 +540,55 @@ impl BucketMetadataSys { let bucket = bucket.clone(); futures.push(async move { sleep(Duration::from_millis(30)).await; - let _ = api - .heal_bucket( - &bucket, - &HealOpts { - recreate: true, - ..Default::default() - }, - ) - .await; - load_bucket_metadata_parse_with_presence(self.api.clone(), bucket.as_str(), true).await + match mode { + MetadataLoadMode::Initial => { + let _ = api + .heal_bucket( + &bucket, + &HealOpts { + recreate: true, + ..Default::default() + }, + ) + .await; + let (bm, persisted) = + load_bucket_metadata_parse_with_presence(self.api.clone(), bucket.as_str(), true).await?; + if persisted { + self.set(bucket, Arc::new(bm)).await; + } else { + let _publish_guard = self.metadata_publish_lock.lock().await; + let mut map = self.metadata_map.write().await; + map.entry(bucket).or_insert_with(|| Arc::new(bm)); + } + } + MetadataLoadMode::Refresh => { + let expected = self.metadata_map.read().await.get(&bucket).cloned(); + let heal_lock = api.new_ns_lock(&bucket, &bucket).await?; + let heal_guard = heal_lock.get_read_lock(crate::set_disk::get_lock_acquire_timeout()).await?; + await_bucket_namespace_operation( + Some(&heal_guard), + &bucket, + "bucket metadata refresh heal", + api.heal_bucket(&bucket, &HealOpts::default()), + ) + .await?; + drop(heal_guard); + let (bm, persisted) = + load_bucket_metadata_parse_with_presence(self.api.clone(), bucket.as_str(), true).await?; + let publish_lock = api.new_ns_lock(&bucket, &bucket).await?; + let guard = publish_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + if guard.is_lock_lost() { + return Err(Error::other(format!( + "bucket namespace lock was lost before bucket metadata refresh publish: {bucket}" + ))); + } + self.publish_refresh_if_unchanged(&bucket, expected.as_ref(), bm, persisted) + .await; + } + } + Ok::<(), Error>(()) }); } @@ -542,26 +596,7 @@ impl BucketMetadataSys { for (idx, res) in results.into_iter().enumerate() { match res { - Ok((bm, persisted)) => { - if let Some(bucket) = buckets.get(idx) { - if persisted { - self.set(bucket.clone(), Arc::new(bm)).await; - } else { - // A fabricated default (no persisted metadata - // readable right now) must never REPLACE an - // existing entry: the periodic refresh would - // otherwise downgrade a lock-enabled bucket to an - // authoritative "no lock" default on a transient - // ConfigNotFound, disabling the object-lock - // delete gate and wiping its target/durability - // sync state. Insert-if-vacant keeps the startup - // behavior for legacy buckets without a metadata - // file, atomically under the map write lock. - let mut map = self.metadata_map.write().await; - map.entry(bucket.clone()).or_insert_with(|| Arc::new(bm)); - } - } - } + Ok(()) => {} Err(e) => { error!("Unable to load bucket metadata, will be retried: {:?}", e); if let Some(bucket) = buckets.get(idx) { @@ -572,6 +607,32 @@ impl BucketMetadataSys { } } + async fn publish_refresh_if_unchanged( + &self, + bucket: &str, + expected: Option<&Arc>, + metadata: BucketMetadata, + persisted: bool, + ) { + if !persisted { + return; + } + let _publish_guard = self.metadata_publish_lock.lock().await; + let metadata = Arc::new(metadata); + let mut map = self.metadata_map.write().await; + let unchanged = expected + .zip(map.get(bucket)) + .is_some_and(|(expected, current)| Arc::ptr_eq(expected, current)); + if !unchanged { + return; + } + map.insert(bucket.to_string(), Arc::clone(&metadata)); + drop(map); + self.absent_metadata.invalidate(bucket).await; + sync_bucket_target_sys(bucket, &metadata).await; + sync_bucket_durability(bucket, &metadata); + } + pub async fn get(&self, bucket: &str) -> Result> { if is_meta_bucketname(bucket) { return Err(Error::ConfigNotFound); @@ -587,6 +648,7 @@ impl BucketMetadataSys { pub async fn set(&self, bucket: String, bm: Arc) { if !is_meta_bucketname(&bucket) { + let _publish_guard = self.metadata_publish_lock.lock().await; let mut map = self.metadata_map.write().await; map.insert(bucket.clone(), bm.clone()); drop(map); @@ -604,6 +666,7 @@ impl BucketMetadataSys { if is_meta_bucketname(bucket) { return false; } + let _publish_guard = self.metadata_publish_lock.lock().await; let mut map = self.metadata_map.write().await; let removed = map.remove(bucket).is_some(); drop(map); @@ -735,7 +798,24 @@ impl BucketMetadataSys { return Ok((Arc::new(bm), true)); } - let (bm, persisted) = match load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true).await { + let lock = self.api.new_ns_lock(bucket, bucket).await?; + let guard = lock.get_read_lock(crate::set_disk::get_lock_acquire_timeout()).await?; + #[cfg(test)] + if self.lazy_load_lock_probe.load(std::sync::atomic::Ordering::Relaxed) { + let competing = self.api.new_ns_lock(bucket, bucket).await?; + assert!( + competing.get_write_lock(Duration::from_millis(20)).await.is_err(), + "lazy metadata IO must start while the bucket namespace read lock is held" + ); + } + let (bm, persisted) = match 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 { @@ -759,9 +839,33 @@ impl BucketMetadataSys { // defaults for buckets listed on disk — legacy buckets without a // metadata file — but never lets one replace an existing entry.) if persisted { + await_bucket_namespace_operation( + Some(&guard), + bucket, + "lazy bucket metadata existence check", + Box::pin(async { + self.api + .peer_sys + .get_bucket_info(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()) + .await + .map(|_| ()) + .map_err(Into::into) + }), + ) + .await?; + if guard.is_lock_lost() { + return Err(Error::other(format!( + "bucket namespace lock was lost before lazy bucket metadata publish: {bucket}" + ))); + } + let _publish_guard = self.metadata_publish_lock.lock().await; let mut map = self.metadata_map.write().await; + if let Some(current) = map.get(bucket) { + return Ok((Arc::clone(current), true)); + } map.insert(bucket.to_string(), bm.clone()); drop(map); + self.absent_metadata.invalidate(bucket).await; sync_bucket_target_sys(bucket, &bm).await; sync_bucket_durability(bucket, &bm); } else { @@ -1047,12 +1151,13 @@ mod tests { /// Pins the fail-closed caching contract of the lazy `get_config` path /// and the refresh no-replace rule: fabricated defaults are returned but /// never served by the map-only `get()`, persisted metadata is cached on - /// lazy load (superseding a recorded absence), and a refresh-load miss - /// never replaces an existing entry. + /// lazy load (superseding a recorded absence), a refresh-load miss never + /// replaces an existing entry or heals a deleted bucket, and initial load + /// still heals buckets discovered from storage. #[tokio::test] async fn get_config_never_caches_fabricated_defaults_as_authoritative() { - let (_dirs, ecstore) = isolated_store_over_temp_disks().await; - let sys = BucketMetadataSys::new(ecstore); + let (dirs, ecstore) = isolated_store_over_temp_disks().await; + let sys = Arc::new(BucketMetadataSys::new(ecstore)); // (a) Miss: the fabricated default is returned but not cached. let (bm, _) = sys @@ -1078,6 +1183,9 @@ mod tests { let mut persisted = BucketMetadata::new("absent-bucket"); persisted.policy_config_json = b"persisted-marker".to_vec(); sys.persist_and_set(persisted).await.expect("metadata should persist"); + for dir in &dirs { + std::fs::create_dir_all(dir.path().join("absent-bucket")).expect("persisted bucket directory should be created"); + } sys.metadata_map.write().await.clear(); let _ = sys .get_config("absent-bucket") @@ -1089,14 +1197,47 @@ mod tests { .expect("lazily loaded persisted metadata must be cached"); assert_eq!(cached.policy_config_json, b"persisted-marker".to_vec()); - // (c) A refresh-load miss (no persisted metadata readable) must not - // replace an existing entry. + // (c) Persisted metadata left behind after physical deletion must not + // be lazily republished as a live bucket generation. + let mut deleted_lazy = BucketMetadata::new("deleted-lazy-bucket"); + deleted_lazy.policy_config_json = b"stale-generation".to_vec(); + sys.persist_and_set(deleted_lazy) + .await + .expect("stale metadata should persist"); + sys.metadata_map.write().await.remove("deleted-lazy-bucket"); + assert!( + sys.get_config("deleted-lazy-bucket").await.is_err(), + "lazy load must fail when the physical bucket no longer exists" + ); + assert!(sys.get("deleted-lazy-bucket").await.is_err()); + + // (d) The namespace generation fence must be acquired before lazy + // metadata IO, so a writer can replace the generation atomically. + let fenced_bucket = "fenced-lazy-bucket"; + for dir in &dirs { + std::fs::create_dir_all(dir.path().join(fenced_bucket)).unwrap(); + } + let mut old_fenced = BucketMetadata::new(fenced_bucket); + old_fenced.policy_config_json = b"old-fenced-generation".to_vec(); + sys.persist_and_set(old_fenced).await.unwrap(); + sys.metadata_map.write().await.remove(fenced_bucket); + sys.lazy_load_lock_probe.store(true, std::sync::atomic::Ordering::Relaxed); + let (loaded, _) = sys.get_config(fenced_bucket).await.unwrap(); + sys.lazy_load_lock_probe.store(false, std::sync::atomic::Ordering::Relaxed); + assert_eq!(loaded.policy_config_json, b"old-fenced-generation".to_vec()); + + // (e) A refresh-load miss for a bucket that still exists must not + // replace an existing entry with a fabricated default. let mut kept = BucketMetadata::new("kept-bucket"); kept.policy_config_json = b"kept-marker".to_vec(); sys.set("kept-bucket".to_string(), Arc::new(kept)).await; + for dir in &dirs { + std::fs::create_dir_all(dir.path().join("kept-bucket")).expect("kept bucket directory should be created"); + } let mut failed = HashSet::new(); let refresh_targets = vec!["kept-bucket".to_string()]; - sys.concurrent_load(&refresh_targets, &mut failed).await; + sys.concurrent_load(&refresh_targets, &mut failed, MetadataLoadMode::Refresh) + .await; let kept = sys .get("kept-bucket") .await @@ -1106,6 +1247,50 @@ mod tests { b"kept-marker".to_vec(), "a fabricated refresh default must not replace real metadata" ); + + // (f) A stale cache entry for a physically deleted bucket must not + // recreate the bucket during periodic refresh. + sys.set("deleted-bucket".to_string(), Arc::new(BucketMetadata::new("deleted-bucket"))) + .await; + let deleted_targets = vec!["deleted-bucket".to_string()]; + sys.concurrent_load(&deleted_targets, &mut failed, MetadataLoadMode::Refresh) + .await; + assert!( + dirs.iter().all(|dir| !dir.path().join("deleted-bucket").exists()), + "periodic refresh must not recreate a bucket from stale cached metadata" + ); + + // (g) Metadata loaded for an old bucket generation must not replace + // metadata published by delete plus same-name recreation. + let old = Arc::new(BucketMetadata::new("recreated-bucket")); + sys.set("recreated-bucket".to_string(), Arc::clone(&old)).await; + let mut recreated = BucketMetadata::new("recreated-bucket"); + recreated.policy_config_json = b"new-generation".to_vec(); + sys.set("recreated-bucket".to_string(), Arc::new(recreated)).await; + let mut stale = BucketMetadata::new("recreated-bucket"); + stale.policy_config_json = b"old-generation".to_vec(); + sys.publish_refresh_if_unchanged("recreated-bucket", Some(&old), stale, true) + .await; + assert_eq!(sys.get("recreated-bucket").await.unwrap().policy_config_json, b"new-generation".to_vec()); + + // (f) Refresh retains periodic healing for a partially missing bucket. + sys.set("partial-bucket".to_string(), Arc::new(BucketMetadata::new("partial-bucket"))) + .await; + for dir in dirs.iter().take(3) { + std::fs::create_dir_all(dir.path().join("partial-bucket")).unwrap(); + } + sys.concurrent_load(&["partial-bucket".to_string()], &mut failed, MetadataLoadMode::Refresh) + .await; + assert!(dirs.iter().all(|dir| dir.path().join("partial-bucket").is_dir())); + + // (g) Initial discovery retains the historical unconditional heal. + let initial_targets = vec!["initial-bucket".to_string()]; + sys.concurrent_load(&initial_targets, &mut failed, MetadataLoadMode::Initial) + .await; + assert!( + dirs.iter().all(|dir| dir.path().join("initial-bucket").is_dir()), + "initial load must heal buckets discovered from storage" + ); } #[tokio::test] @@ -1130,12 +1315,18 @@ mod tests { #[tokio::test] async fn update_config_with_persists_tagging_rewrite_across_disk_reload() { use crate::bucket::metadata::BUCKET_TAGGING_CONFIG; + use crate::storage_api_contracts::bucket::MakeBucketOptions; use s3s::dto::Tag; let (_dirs, ecstore) = isolated_store_over_temp_disks().await; - let mut sys = BucketMetadataSys::new(ecstore); let bucket = "swift-tagging-bucket"; + ecstore + .peer_sys + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket volume should be created"); + let mut sys = BucketMetadataSys::new(ecstore); sys.persist_and_set(BucketMetadata::new(bucket)) .await .expect("initial metadata should persist"); diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index 17d42e84d..67d04ab23 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -89,6 +89,22 @@ fn reduce_pool_write_quorum_errs(per_pool_errs: &[Option]) -> Option]) -> Result<()> { + if opts.recreate { + return Ok(()); + } + if let Some(err) = pool_errs + .iter() + .flatten() + .find(|err| **err != Error::DiskNotFound && **err != Error::VolumeNotFound) + { + return Err(err.clone()); + } + opts.remove = is_all_buckets_not_found(pool_errs); + opts.recreate = !opts.remove; + Ok(()) +} + #[async_trait] pub trait PeerS3Client: Debug + Sync + Send + 'static { async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result; @@ -159,10 +175,7 @@ impl S3PeerSys { pool_errs.push(reduce_pool_write_quorum_errs(&per_pool_errs)); } - if !opts.recreate { - opts.remove = is_all_buckets_not_found(&pool_errs); - opts.recreate = !opts.remove; - } + resolve_heal_bucket_mode(&mut opts, &pool_errs)?; let mut futures = Vec::new(); let heal_bucket_results = Arc::new(RwLock::new(vec![HealResultItem::default(); self.clients.len()])); @@ -1618,6 +1631,30 @@ mod tests { assert_eq!(err, Error::VolumeExists); } + #[test] + fn heal_bucket_mode_fails_closed_on_incomplete_topology() { + let mut opts = HealOpts::default(); + assert_eq!( + resolve_heal_bucket_mode(&mut opts, &[Some(Error::ErasureWriteQuorum)]), + Err(Error::ErasureWriteQuorum) + ); + assert!(!opts.recreate); + assert!(!opts.remove); + } + + #[test] + fn heal_bucket_mode_distinguishes_deleted_and_partial_buckets() { + let mut deleted = HealOpts::default(); + resolve_heal_bucket_mode(&mut deleted, &[Some(Error::VolumeNotFound)]).unwrap(); + assert!(deleted.remove); + assert!(!deleted.recreate); + + let mut partial = HealOpts::default(); + resolve_heal_bucket_mode(&mut partial, &[None, Some(Error::VolumeNotFound)]).unwrap(); + assert!(!partial.remove); + assert!(partial.recreate); + } + #[tokio::test] async fn test_make_bucket_reduces_quorum_by_pool_participants() { let peer_sys = S3PeerSys { diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index fcdc67a44..4d5d4840b 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -87,7 +87,7 @@ fn bucket_deleted_marker_volume(bucket: &str) -> String { format!("{RUSTFS_META_BUCKET}/{}", bucket_deleted_marker_prefix(bucket)) } -async fn await_bucket_namespace_operation( +pub(crate) async fn await_bucket_namespace_operation( guard: Option<&rustfs_lock::NamespaceLockGuard>, bucket: &str, operation: &'static str, diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index bad794f01..4c9ea3039 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -141,6 +141,7 @@ fn should_enqueue_transition_immediately(oi: &ObjectInfo) -> bool { const MAX_UPLOADS_LIST: usize = 10000; mod bucket; +pub(crate) use bucket::await_bucket_namespace_operation; mod heal; mod heal_walk; pub use heal_walk::HealWalkVersion;