diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 70a84851f..9896b0ea7 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -35,7 +35,10 @@ use s3s::dto::{ }; use std::collections::HashSet; use std::time::Duration; -use std::{collections::HashMap, sync::Arc}; +use std::{ + collections::HashMap, + sync::{Arc, Mutex as StdMutex, Weak}, +}; use time::OffsetDateTime; use tokio::sync::{Mutex, RwLock}; use tokio::time::sleep; @@ -91,18 +94,17 @@ pub async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) -> Result<( Ok(()) } -/// Peer LoadBucketMetadata entry point; see -/// [`BucketMetadataSys::reload_from_store`] for the caching contract. -/// -/// The outer write guard spans the disk load, mirroring [`update`]: every -/// other cache installer holds this lock (read or write), so the snapshot -/// read here can never land after — and roll back — a newer concurrent -/// install, and the install-plus-registry-sync sequence stays atomic -/// against concurrent removes and reloads. -pub async fn reload_bucket_metadata(bucket: &str) -> Result<()> { - let sys = get_bucket_metadata_sys()?; - let lock = sys.write().await; - lock.reload_from_store(bucket).await +pub async fn reload_bucket_metadata(api: Arc, bucket: &str) -> Result<()> { + if is_meta_bucketname(bucket) { + return Err(Error::other("errInvalidArgument")); + } + let namespace_lock = api.new_ns_lock(bucket, bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + let sys = bucket_metadata_sys_of(&api.ctx)?; + let lock = sys.read().await; + lock.reload_from_store_under_namespace(bucket, &namespace_guard).await } /// Drop a bucket's cached metadata from the in-memory map. @@ -159,9 +161,7 @@ async fn refresh_buckets_metadata_once(sys: Arc>) { let mut failed_buckets = HashSet::new(); for chunk in buckets.chunks(count) { - let sys = sys.read().await; - sys.concurrent_load(chunk, &mut failed_buckets, MetadataLoadMode::Refresh) - .await; + BucketMetadataSys::concurrent_refresh_load(Arc::clone(&sys), chunk, &mut failed_buckets).await; } if !failed_buckets.is_empty() { @@ -478,11 +478,42 @@ pub async fn list_bucket_targets(bucket: &str) -> Result { /// notification was lost; the capacity bounds memory under bogus-name floods. const ABSENT_BUCKET_METADATA_TTL: Duration = Duration::from_secs(30); const ABSENT_BUCKET_METADATA_MAX_ENTRIES: u64 = 10_000; +const PEER_METADATA_NOT_PERSISTED: &str = "no persisted bucket metadata readable; peer cache left unchanged"; +#[derive(Debug)] +struct MetadataPublishLockRegistry { + locks: StdMutex>>>, +} + +#[derive(Debug)] +struct MetadataPublishLockState { + bucket: String, + registry: Weak, + lock: Weak>, +} + +impl Drop for MetadataPublishLockState { + fn drop(&mut self) { + let Some(registry) = self.registry.upgrade() else { + return; + }; + let mut locks = registry.locks.lock().unwrap_or_else(|poisoned| poisoned.into_inner()); + if locks.get(&self.bucket).is_some_and(|current| current.ptr_eq(&self.lock)) { + locks.remove(&self.bucket); + } + } +} + +#[derive(Debug)] +struct MetadataPublishGuard { + _guard: tokio::sync::OwnedMutexGuard, +} #[derive(Debug)] pub struct BucketMetadataSys { metadata_map: RwLock>>, - metadata_publish_lock: Mutex<()>, + /// Serializes metadata-map commits and their derived cache updates for one + /// bucket. Namespace locks, when present, are acquired before this lock. + metadata_publish_locks: Arc, #[cfg(test)] lazy_load_lock_probe: std::sync::atomic::AtomicBool, /// Buckets recently observed to have no persisted metadata. Serving the @@ -500,7 +531,9 @@ impl BucketMetadataSys { pub fn new(api: Arc) -> Self { Self { metadata_map: RwLock::new(HashMap::new()), - metadata_publish_lock: Mutex::new(()), + metadata_publish_locks: Arc::new(MetadataPublishLockRegistry { + locks: StdMutex::new(HashMap::new()), + }), #[cfg(test)] lazy_load_lock_probe: std::sync::atomic::AtomicBool::new(false), absent_metadata: moka::future::Cache::builder() @@ -516,6 +549,65 @@ impl BucketMetadataSys { self.api.clone() } + fn metadata_publish_lock(&self, bucket: &str) -> Arc> { + let mut locks = self + .metadata_publish_locks + .locks + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + locks.get(bucket).and_then(Weak::upgrade).unwrap_or_else(|| { + let lock = Arc::new_cyclic(|lock| { + Mutex::new(MetadataPublishLockState { + bucket: bucket.to_string(), + registry: Arc::downgrade(&self.metadata_publish_locks), + lock: lock.clone(), + }) + }); + locks.insert(bucket.to_string(), Arc::downgrade(&lock)); + lock + }) + } + + async fn lock_metadata_publish( + &self, + bucket: &str, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + operation: &'static str, + ) -> Result { + let lock = self.metadata_publish_lock(bucket); + let guard = await_bucket_namespace_operation(Some(namespace_guard), bucket, operation, async { + Ok(MetadataPublishGuard { + _guard: lock.lock_owned().await, + }) + }) + .await?; + if namespace_guard.is_lock_lost() { + return Err(Error::other(format!("bucket namespace lock was lost before {operation}: {bucket}"))); + } + Ok(guard) + } + + async fn bucket_exists( + &self, + bucket: &str, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + operation: &'static str, + ) -> Result { + 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()) + .await + { + Ok(_) => Ok(true), + Err(crate::disk::error::Error::VolumeNotFound) => Ok(false), + Err(err) => Err(err.into()), + } + }) + .await + } + pub async fn init(&mut self, buckets: Vec) { let _ = self.init_internal(buckets).await; } @@ -554,55 +646,16 @@ impl BucketMetadataSys { let bucket = bucket.clone(); futures.push(async move { sleep(Duration::from_millis(30)).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>(()) + let expected = match mode { + MetadataLoadMode::Initial => None, + MetadataLoadMode::Refresh => self.metadata_map.read().await.get(&bucket).cloned(), + }; + let namespace_lock = api.new_ns_lock(&bucket, &bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + self.load_bucket_under_namespace(&bucket, mode, expected.as_ref(), &namespace_guard) + .await }); } @@ -621,30 +674,134 @@ impl BucketMetadataSys { } } - async fn publish_refresh_if_unchanged( + async fn concurrent_refresh_load(sys: Arc>, buckets: &[String], failed_buckets: &mut HashSet) { + let mut futures = Vec::with_capacity(buckets.len()); + for bucket in buckets { + let sys = Arc::clone(&sys); + let bucket = bucket.clone(); + futures.push(async move { + sleep(Duration::from_millis(30)).await; + let api = sys.read().await.api.clone(); + let namespace_lock = api.new_ns_lock(&bucket, &bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + let metadata_sys = sys.read().await; + let expected = metadata_sys.metadata_map.read().await.get(&bucket).cloned(); + metadata_sys + .load_bucket_under_namespace(&bucket, MetadataLoadMode::Refresh, expected.as_ref(), &namespace_guard) + .await + }); + } + let results = join_all(futures).await; + for (idx, result) in results.into_iter().enumerate() { + if let Err(err) = result { + error!("Unable to load bucket metadata, will be retried: {:?}", err); + if let Some(bucket) = buckets.get(idx) { + failed_buckets.insert(bucket.clone()); + } + } + } + } + + async fn load_bucket_under_namespace( + &self, + bucket: &str, + mode: MetadataLoadMode, + 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? + { + if matches!(mode, MetadataLoadMode::Refresh) { + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "stale bucket metadata removal") + .await?; + let removed = self.metadata_map.write().await.remove(bucket).is_some(); + if removed { + BucketTargetSys::get().delete(bucket).await; + clear_bucket_durability(bucket); + } + } + return Ok(()); + } + + let (bm, persisted) = await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "bucket metadata load", + load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true), + ) + .await?; + match mode { + MetadataLoadMode::Initial if persisted => { + let bm = Arc::new(bm); + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "initial bucket metadata publish") + .await?; + self.metadata_map.write().await.insert(bucket.to_string(), Arc::clone(&bm)); + self.absent_metadata.invalidate(bucket).await; + sync_bucket_target_sys(bucket, &bm).await; + sync_bucket_durability(bucket, &bm); + } + MetadataLoadMode::Initial => { + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "initial bucket metadata publish") + .await?; + self.metadata_map + .write() + .await + .entry(bucket.to_string()) + .or_insert_with(|| Arc::new(bm)); + } + MetadataLoadMode::Refresh => { + self.publish_if_unchanged(bucket, expected, bm, persisted, namespace_guard) + .await?; + } + } + Ok(()) + } + + async fn publish_if_unchanged( &self, bucket: &str, expected: Option<&Arc>, metadata: BucketMetadata, persisted: bool, - ) { + namespace_guard: &rustfs_lock::NamespaceLockGuard, + ) -> Result<()> { if !persisted { - return; + return Ok(()); } - let _publish_guard = self.metadata_publish_lock.lock().await; + let _publish_guard = self + .lock_metadata_publish(bucket, namespace_guard, "refreshed bucket metadata publish") + .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)); + let unchanged = match (expected, map.get(bucket)) { + (None, None) => true, + (Some(expected), Some(current)) => Arc::ptr_eq(expected, current), + _ => false, + }; if !unchanged { - return; + return Ok(()); } 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); + Ok(()) } pub async fn get(&self, bucket: &str) -> Result> { @@ -662,7 +819,8 @@ 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 publish_lock = self.metadata_publish_lock(&bucket); + let _publish_guard = publish_lock.lock().await; let mut map = self.metadata_map.write().await; map.insert(bucket.clone(), bm.clone()); drop(map); @@ -673,43 +831,6 @@ impl BucketMetadataSys { } } - /// Reload `bucket`'s metadata from this system's own store and cache it, - /// refusing to treat a load miss as authoritative (the peer - /// LoadBucketMetadata notification path, [`reload_bucket_metadata`]). - /// - /// Only metadata actually read from persisted storage reaches the cache. - /// On a miss the fabricated default is discarded and an error is - /// returned: installing it would let a transient ConfigNotFound during - /// the notification overwrite a lock-enabled bucket's cached metadata - /// with an authoritative "no Object Lock" default, disabling the - /// batch-delete retention gate (`object_lock_delete_check_required`) on - /// this node until the next refresh. A miss is also not treated as - /// deletion: bucket deletion propagates through the dedicated - /// DeleteBucketMetadata notification ([`remove_bucket_metadata`]), which - /// is best-effort — a reload racing it can still re-install a just - /// deleted bucket's entry (pre-existing, bounded by the next delete or - /// restart) — but a reload miss removing entries would turn every - /// transient quorum dip into dropped metadata and spurious - /// target/durability teardown. - /// - /// The peer-visible error text is deliberately fixed: the notifying peer - /// matches error strings against network-failure needles - /// (`is_network_like_error`), so interpolating a caller-controlled - /// bucket name here could mark a healthy peer offline. - /// - /// Lock order: the caller holds the outer metadata-sys guard, and the - /// load acquires the namespace lock on the bucket's metadata config - /// object — the same `outer guard → meta-config namespace lock` order - /// `update`'s load takes; no path acquires these in reverse. - pub(crate) async fn reload_from_store(&self, bucket: &str) -> Result<()> { - let (bm, persisted) = load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true).await?; - if !persisted { - return Err(Error::other("no persisted bucket metadata readable; peer cache left unchanged")); - } - self.set(bucket.to_string(), Arc::new(bm)).await; - Ok(()) - } - /// Remove a bucket's cached metadata from the in-memory map. /// /// Returns `true` if an entry was present. Reserved meta buckets are ignored. @@ -717,7 +838,8 @@ impl BucketMetadataSys { if is_meta_bucketname(bucket) { return false; } - let _publish_guard = self.metadata_publish_lock.lock().await; + let publish_lock = self.metadata_publish_lock(bucket); + let _publish_guard = publish_lock.lock().await; let mut map = self.metadata_map.write().await; let removed = map.remove(bucket).is_some(); drop(map); @@ -831,6 +953,49 @@ impl BucketMetadataSys { load_bucket_metadata(self.api.clone(), bucket).await } + /// Reload persisted metadata under the bucket namespace generation fence. + /// + /// A miss is never published as an authoritative default, and a snapshot + /// read before delete plus same-name recreation cannot replace the new + /// generation. + pub(crate) async fn reload_from_store(&self, bucket: &str) -> Result<()> { + if is_meta_bucketname(bucket) { + return Err(Error::other("errInvalidArgument")); + } + + let namespace_lock = self.api.new_ns_lock(bucket, bucket).await?; + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await?; + self.reload_from_store_under_namespace(bucket, &namespace_guard).await + } + + async fn reload_from_store_under_namespace( + &self, + bucket: &str, + namespace_guard: &rustfs_lock::NamespaceLockGuard, + ) -> Result<()> { + let expected = self.metadata_map.read().await.get(bucket).cloned(); + if !self + .bucket_exists(bucket, namespace_guard, "peer bucket metadata existence check") + .await? + { + return Err(Error::other(PEER_METADATA_NOT_PERSISTED)); + } + let (metadata, persisted) = await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "peer bucket metadata load", + load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true), + ) + .await?; + if !persisted { + return Err(Error::other(PEER_METADATA_NOT_PERSISTED)); + } + self.publish_if_unchanged(bucket, expected.as_ref(), metadata, true, namespace_guard) + .await + } + pub async fn get_config(&self, bucket: &str) -> Result<(Arc, bool)> { let has_bm = { let map = self.metadata_map.read().await; @@ -909,7 +1074,9 @@ impl BucketMetadataSys { "bucket namespace lock was lost before lazy bucket metadata publish: {bucket}" ))); } - let _publish_guard = self.metadata_publish_lock.lock().await; + let _publish_guard = self + .lock_metadata_publish(bucket, &guard, "lazy bucket metadata publish") + .await?; let mut map = self.metadata_map.write().await; if let Some(current) = map.get(bucket) { return Ok((Arc::clone(current), true)); @@ -1247,6 +1414,11 @@ mod tests { .await .expect("lazily loaded persisted metadata must be cached"); assert_eq!(cached.policy_config_json, b"persisted-marker".to_vec()); + sys.metadata_map.write().await.clear(); + sys.reload_from_store("absent-bucket") + .await + .expect("peer reload should publish persisted metadata into a cold cache"); + assert_eq!(sys.get("absent-bucket").await.unwrap().policy_config_json, b"persisted-marker".to_vec()); // (c) Persisted metadata left behind after physical deletion must not // be lazily republished as a live bucket generation. @@ -1299,8 +1471,8 @@ mod tests { "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. + // (f) A stale cache entry for a physically deleted bucket must be + // removed without recreating 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()]; @@ -1310,8 +1482,30 @@ mod tests { dirs.iter().all(|dir| !dir.path().join("deleted-bucket").exists()), "periodic refresh must not recreate a bucket from stale cached metadata" ); + assert!( + sys.get("deleted-bucket").await.is_err(), + "periodic refresh must remove stale cached metadata" + ); - // (g) Metadata loaded for an old bucket generation must not replace + // (g) Persisted metadata left behind after physical deletion must not + // keep the deleted generation authoritative during refresh. + let mut deleted_persisted = BucketMetadata::new("deleted-persisted-bucket"); + deleted_persisted.policy_config_json = b"stale-persisted-generation".to_vec(); + sys.persist_and_set(deleted_persisted) + .await + .expect("stale metadata should persist"); + sys.concurrent_load(&["deleted-persisted-bucket".to_string()], &mut failed, MetadataLoadMode::Refresh) + .await; + assert!( + sys.get("deleted-persisted-bucket").await.is_err(), + "refresh must remove persisted metadata for a physically absent bucket" + ); + assert!( + sys.reload_from_store("deleted-persisted-bucket").await.is_err(), + "peer reload must not publish stale metadata for an absent bucket" + ); + + // (h) 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; @@ -1320,11 +1514,27 @@ mod tests { 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()); + let namespace_lock = sys + .api + .new_ns_lock("recreated-bucket", "recreated-bucket") + .await + .expect("namespace lock should be created"); + let namespace_guard = namespace_lock + .get_read_lock(crate::set_disk::get_lock_acquire_timeout()) + .await + .expect("namespace read lock should be acquired"); + sys.publish_if_unchanged("recreated-bucket", Some(&old), stale, true, &namespace_guard) + .await + .expect("stale refresh publish should be fenced"); + assert_eq!( + sys.get("recreated-bucket") + .await + .expect("recreated bucket metadata should remain cached") + .policy_config_json, + b"new-generation".to_vec() + ); - // (f) Refresh retains periodic healing for a partially missing bucket. + // (i) 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) { @@ -1334,7 +1544,22 @@ mod tests { .await; assert!(dirs.iter().all(|dir| dir.path().join("partial-bucket").is_dir())); - // (g) Initial discovery retains the historical unconditional heal. + // (j) A stale initial snapshot must not recreate a bucket that has + // disappeared from every disk. + let stale_initial_targets = vec!["deleted-initial-bucket".to_string()]; + sys.concurrent_load(&stale_initial_targets, &mut failed, MetadataLoadMode::Initial) + .await; + assert!( + dirs.iter().all(|dir| !dir.path().join("deleted-initial-bucket").exists()), + "initial load must not recreate a bucket absent from every disk" + ); + + // (k) Initial discovery still heals a bucket present on part of the + // storage topology. + for dir in dirs.iter().take(3) { + std::fs::create_dir_all(dir.path().join("initial-bucket")) + .expect("partial initial bucket directory should be created"); + } let initial_targets = vec!["initial-bucket".to_string()]; sys.concurrent_load(&initial_targets, &mut failed, MetadataLoadMode::Initial) .await; @@ -1344,6 +1569,46 @@ mod tests { ); } + #[tokio::test] + async fn metadata_publish_locks_are_isolated_per_bucket() { + let (_dirs, ecstore) = isolated_store_over_temp_disks().await; + let sys = Arc::new(BucketMetadataSys::new(ecstore)); + let first_lock = sys.metadata_publish_lock("blocked-bucket"); + let same_lock = sys.metadata_publish_lock("blocked-bucket"); + assert!(Arc::ptr_eq(&first_lock, &same_lock)); + let first_guard = first_lock.lock().await; + let cancelled_waiter_lock = sys.metadata_publish_lock("blocked-bucket"); + let cancelled_waiter = tokio::spawn(async move { + let _guard = cancelled_waiter_lock.lock_owned().await; + }); + tokio::task::yield_now().await; + cancelled_waiter.abort(); + assert!(cancelled_waiter.await.unwrap_err().is_cancelled()); + let other_bucket = "other-bucket".to_string(); + let other_lock = sys.metadata_publish_lock(&other_bucket); + assert!(!Arc::ptr_eq(&first_lock, &other_lock)); + + timeout( + Duration::from_secs(1), + sys.set(other_bucket.clone(), Arc::new(BucketMetadata::new(&other_bucket))), + ) + .await + .expect("one bucket publish lock must not block another bucket"); + assert!(sys.get(&other_bucket).await.is_ok()); + + drop(first_guard); + drop(first_lock); + drop(same_lock); + drop(other_lock); + assert!( + sys.metadata_publish_locks + .locks + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .is_empty() + ); + } + #[tokio::test] async fn get_bucket_policy_rejects_malformed_cached_policy() { let (_dirs, ecstore) = isolated_store_over_temp_disks().await; @@ -1492,7 +1757,7 @@ mod tests { /// default and disable the batch-delete retention gate on this peer. #[tokio::test] async fn peer_reload_never_caches_fabricated_defaults_as_authoritative() { - let (_dirs, ecstore) = isolated_store_over_temp_disks().await; + let (dirs, ecstore) = isolated_store_over_temp_disks().await; let sys = BucketMetadataSys::new(ecstore.clone()); // (a) Miss with no cached entry: the reload fails and installs nothing. @@ -1530,6 +1795,9 @@ mod tests { let mut persisted = BucketMetadata::new("reload-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("reload-bucket")).expect("physical bucket should exist before reload"); + } let mut stale = BucketMetadata::new("reload-bucket"); stale.policy_config_json = b"stale-cache-marker".to_vec(); sys.set("reload-bucket".to_string(), Arc::new(stale)).await; diff --git a/rustfs/src/storage/rpc/node_service/bucket.rs b/rustfs/src/storage/rpc/node_service/bucket.rs index 2f9c490c7..091fe028f 100644 --- a/rustfs/src/storage/rpc/node_service/bucket.rs +++ b/rustfs/src/storage/rpc/node_service/bucket.rs @@ -66,14 +66,14 @@ impl NodeService { })); } - let Some(_store) = self.resolve_object_store() else { + let Some(store) = self.resolve_object_store() else { return Ok(Response::new(LoadBucketMetadataResponse { success: false, error_info: Some("errServerNotInitialized".to_string()), })); }; - match reload_bucket_metadata(&bucket).await { + match reload_bucket_metadata(store, &bucket).await { Ok(()) => { if scanner_maintenance_change { rustfs_scanner::record_scanner_maintenance_change(&bucket); diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 1b9b27056..9c891d9f6 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -1442,8 +1442,8 @@ pub(crate) async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) -> R ecstore_bucket::metadata_sys::set_bucket_metadata(bucket, bm).await } -pub(crate) async fn reload_bucket_metadata(bucket: &str) -> Result<()> { - ecstore_bucket::metadata_sys::reload_bucket_metadata(bucket).await +pub(crate) async fn reload_bucket_metadata(api: Arc, bucket: &str) -> Result<()> { + ecstore_bucket::metadata_sys::reload_bucket_metadata(api, bucket).await } pub(crate) async fn remove_bucket_metadata(bucket: &str) -> Result {