mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-30 01:58:59 +00:00
fix(ecstore): fence bucket metadata generations (#5427)
This commit is contained in:
@@ -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<ECStore>, 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<RwLock<BucketMetadataSys>>) {
|
||||
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<BucketTargets> {
|
||||
/// 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<HashMap<String, Weak<Mutex<MetadataPublishLockState>>>>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct MetadataPublishLockState {
|
||||
bucket: String,
|
||||
registry: Weak<MetadataPublishLockRegistry>,
|
||||
lock: Weak<Mutex<MetadataPublishLockState>>,
|
||||
}
|
||||
|
||||
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<MetadataPublishLockState>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct BucketMetadataSys {
|
||||
metadata_map: RwLock<HashMap<String, Arc<BucketMetadata>>>,
|
||||
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<MetadataPublishLockRegistry>,
|
||||
#[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<ECStore>) -> 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<Mutex<MetadataPublishLockState>> {
|
||||
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<MetadataPublishGuard> {
|
||||
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<bool> {
|
||||
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<String>) {
|
||||
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<RwLock<Self>>, buckets: &[String], failed_buckets: &mut HashSet<String>) {
|
||||
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<BucketMetadata>>,
|
||||
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<BucketMetadata>>,
|
||||
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<Arc<BucketMetadata>> {
|
||||
@@ -662,7 +819,8 @@ impl BucketMetadataSys {
|
||||
|
||||
pub async fn set(&self, bucket: String, bm: Arc<BucketMetadata>) {
|
||||
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<BucketMetadata>, 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;
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<ECStore>, bucket: &str) -> Result<()> {
|
||||
ecstore_bucket::metadata_sys::reload_bucket_metadata(api, bucket).await
|
||||
}
|
||||
|
||||
pub(crate) async fn remove_bucket_metadata(bucket: &str) -> Result<bool> {
|
||||
|
||||
Reference in New Issue
Block a user