mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-31 09:18:28 +00:00
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 <housemecn@gmail.com>
This commit is contained in:
@@ -23,7 +23,7 @@ use crate::error::{Error, Result, is_err_bucket_not_found};
|
|||||||
use crate::runtime::sources as runtime_sources;
|
use crate::runtime::sources as runtime_sources;
|
||||||
use crate::storage_api_contracts::heal::HealOperations as _;
|
use crate::storage_api_contracts::heal::HealOperations as _;
|
||||||
use crate::storage_api_contracts::namespace::NamespaceLocking 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 futures::future::join_all;
|
||||||
use rustfs_common::heal_channel::HealOpts;
|
use rustfs_common::heal_channel::HealOpts;
|
||||||
use rustfs_policy::policy::BucketPolicy;
|
use rustfs_policy::policy::BucketPolicy;
|
||||||
@@ -37,13 +37,19 @@ use std::collections::HashSet;
|
|||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use std::{collections::HashMap, sync::Arc};
|
use std::{collections::HashMap, sync::Arc};
|
||||||
use time::OffsetDateTime;
|
use time::OffsetDateTime;
|
||||||
use tokio::sync::RwLock;
|
use tokio::sync::{Mutex, RwLock};
|
||||||
use tokio::time::sleep;
|
use tokio::time::sleep;
|
||||||
use tokio_util::sync::CancellationToken;
|
use tokio_util::sync::CancellationToken;
|
||||||
use tracing::{error, warn};
|
use tracing::{error, warn};
|
||||||
|
|
||||||
const BUCKET_METADATA_REFRESH_INTERVAL: Duration = Duration::from_secs(15 * 60);
|
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<ECStore>, buckets: Vec<String>) {
|
pub async fn init_bucket_metadata_sys(api: Arc<ECStore>, buckets: Vec<String>) {
|
||||||
// The metadata system is inherently per-store (it holds the store handle
|
// 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
|
// 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<RwLock<BucketMetadataSys>>) {
|
|||||||
|
|
||||||
for chunk in buckets.chunks(count) {
|
for chunk in buckets.chunks(count) {
|
||||||
let sys = sys.read().await;
|
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() {
|
if !failed_buckets.is_empty() {
|
||||||
@@ -461,6 +468,9 @@ const ABSENT_BUCKET_METADATA_MAX_ENTRIES: u64 = 10_000;
|
|||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct BucketMetadataSys {
|
pub struct BucketMetadataSys {
|
||||||
metadata_map: RwLock<HashMap<String, Arc<BucketMetadata>>>,
|
metadata_map: RwLock<HashMap<String, Arc<BucketMetadata>>>,
|
||||||
|
metadata_publish_lock: Mutex<()>,
|
||||||
|
#[cfg(test)]
|
||||||
|
lazy_load_lock_probe: std::sync::atomic::AtomicBool,
|
||||||
/// Buckets recently observed to have no persisted metadata. Serving the
|
/// Buckets recently observed to have no persisted metadata. Serving the
|
||||||
/// fabricated default from here (instead of re-reading disk) keeps the
|
/// fabricated default from here (instead of re-reading disk) keeps the
|
||||||
/// per-request cost of repeated lookups for such names bounded — without
|
/// per-request cost of repeated lookups for such names bounded — without
|
||||||
@@ -476,6 +486,9 @@ impl BucketMetadataSys {
|
|||||||
pub fn new(api: Arc<ECStore>) -> Self {
|
pub fn new(api: Arc<ECStore>) -> Self {
|
||||||
Self {
|
Self {
|
||||||
metadata_map: RwLock::new(HashMap::new()),
|
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()
|
absent_metadata: moka::future::Cache::builder()
|
||||||
.max_capacity(ABSENT_BUCKET_METADATA_MAX_ENTRIES)
|
.max_capacity(ABSENT_BUCKET_METADATA_MAX_ENTRIES)
|
||||||
.time_to_live(ABSENT_BUCKET_METADATA_TTL)
|
.time_to_live(ABSENT_BUCKET_METADATA_TTL)
|
||||||
@@ -502,11 +515,13 @@ impl BucketMetadataSys {
|
|||||||
|
|
||||||
loop {
|
loop {
|
||||||
if buckets.len() < count {
|
if buckets.len() < count {
|
||||||
self.concurrent_load(buckets, &mut failed_buckets).await;
|
self.concurrent_load(buckets, &mut failed_buckets, MetadataLoadMode::Initial)
|
||||||
|
.await;
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
self.concurrent_load(&buckets[..count], &mut failed_buckets).await;
|
self.concurrent_load(&buckets[..count], &mut failed_buckets, MetadataLoadMode::Initial)
|
||||||
|
.await;
|
||||||
|
|
||||||
buckets = &buckets[count..]
|
buckets = &buckets[count..]
|
||||||
}
|
}
|
||||||
@@ -517,7 +532,7 @@ impl BucketMetadataSys {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn concurrent_load(&self, buckets: &[String], failed_buckets: &mut HashSet<String>) {
|
async fn concurrent_load(&self, buckets: &[String], failed_buckets: &mut HashSet<String>, mode: MetadataLoadMode) {
|
||||||
let mut futures = Vec::new();
|
let mut futures = Vec::new();
|
||||||
|
|
||||||
for bucket in buckets.iter() {
|
for bucket in buckets.iter() {
|
||||||
@@ -525,16 +540,55 @@ impl BucketMetadataSys {
|
|||||||
let bucket = bucket.clone();
|
let bucket = bucket.clone();
|
||||||
futures.push(async move {
|
futures.push(async move {
|
||||||
sleep(Duration::from_millis(30)).await;
|
sleep(Duration::from_millis(30)).await;
|
||||||
let _ = api
|
match mode {
|
||||||
.heal_bucket(
|
MetadataLoadMode::Initial => {
|
||||||
&bucket,
|
let _ = api
|
||||||
&HealOpts {
|
.heal_bucket(
|
||||||
recreate: true,
|
&bucket,
|
||||||
..Default::default()
|
&HealOpts {
|
||||||
},
|
recreate: true,
|
||||||
)
|
..Default::default()
|
||||||
.await;
|
},
|
||||||
load_bucket_metadata_parse_with_presence(self.api.clone(), bucket.as_str(), true).await
|
)
|
||||||
|
.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() {
|
for (idx, res) in results.into_iter().enumerate() {
|
||||||
match res {
|
match res {
|
||||||
Ok((bm, persisted)) => {
|
Ok(()) => {}
|
||||||
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));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
error!("Unable to load bucket metadata, will be retried: {:?}", e);
|
error!("Unable to load bucket metadata, will be retried: {:?}", e);
|
||||||
if let Some(bucket) = buckets.get(idx) {
|
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<BucketMetadata>>,
|
||||||
|
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<Arc<BucketMetadata>> {
|
pub async fn get(&self, bucket: &str) -> Result<Arc<BucketMetadata>> {
|
||||||
if is_meta_bucketname(bucket) {
|
if is_meta_bucketname(bucket) {
|
||||||
return Err(Error::ConfigNotFound);
|
return Err(Error::ConfigNotFound);
|
||||||
@@ -587,6 +648,7 @@ impl BucketMetadataSys {
|
|||||||
|
|
||||||
pub async fn set(&self, bucket: String, bm: Arc<BucketMetadata>) {
|
pub async fn set(&self, bucket: String, bm: Arc<BucketMetadata>) {
|
||||||
if !is_meta_bucketname(&bucket) {
|
if !is_meta_bucketname(&bucket) {
|
||||||
|
let _publish_guard = self.metadata_publish_lock.lock().await;
|
||||||
let mut map = self.metadata_map.write().await;
|
let mut map = self.metadata_map.write().await;
|
||||||
map.insert(bucket.clone(), bm.clone());
|
map.insert(bucket.clone(), bm.clone());
|
||||||
drop(map);
|
drop(map);
|
||||||
@@ -604,6 +666,7 @@ impl BucketMetadataSys {
|
|||||||
if is_meta_bucketname(bucket) {
|
if is_meta_bucketname(bucket) {
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
let _publish_guard = self.metadata_publish_lock.lock().await;
|
||||||
let mut map = self.metadata_map.write().await;
|
let mut map = self.metadata_map.write().await;
|
||||||
let removed = map.remove(bucket).is_some();
|
let removed = map.remove(bucket).is_some();
|
||||||
drop(map);
|
drop(map);
|
||||||
@@ -735,7 +798,24 @@ impl BucketMetadataSys {
|
|||||||
return Ok((Arc::new(bm), true));
|
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,
|
Ok(res) => res,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
return if *self.initialized.read().await {
|
return if *self.initialized.read().await {
|
||||||
@@ -759,9 +839,33 @@ impl BucketMetadataSys {
|
|||||||
// defaults for buckets listed on disk — legacy buckets without a
|
// defaults for buckets listed on disk — legacy buckets without a
|
||||||
// metadata file — but never lets one replace an existing entry.)
|
// metadata file — but never lets one replace an existing entry.)
|
||||||
if persisted {
|
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;
|
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());
|
map.insert(bucket.to_string(), bm.clone());
|
||||||
drop(map);
|
drop(map);
|
||||||
|
self.absent_metadata.invalidate(bucket).await;
|
||||||
sync_bucket_target_sys(bucket, &bm).await;
|
sync_bucket_target_sys(bucket, &bm).await;
|
||||||
sync_bucket_durability(bucket, &bm);
|
sync_bucket_durability(bucket, &bm);
|
||||||
} else {
|
} else {
|
||||||
@@ -1047,12 +1151,13 @@ mod tests {
|
|||||||
/// Pins the fail-closed caching contract of the lazy `get_config` path
|
/// Pins the fail-closed caching contract of the lazy `get_config` path
|
||||||
/// and the refresh no-replace rule: fabricated defaults are returned but
|
/// and the refresh no-replace rule: fabricated defaults are returned but
|
||||||
/// never served by the map-only `get()`, persisted metadata is cached on
|
/// never served by the map-only `get()`, persisted metadata is cached on
|
||||||
/// lazy load (superseding a recorded absence), and a refresh-load miss
|
/// lazy load (superseding a recorded absence), a refresh-load miss never
|
||||||
/// never replaces an existing entry.
|
/// replaces an existing entry or heals a deleted bucket, and initial load
|
||||||
|
/// still heals buckets discovered from storage.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn get_config_never_caches_fabricated_defaults_as_authoritative() {
|
async fn get_config_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);
|
let sys = Arc::new(BucketMetadataSys::new(ecstore));
|
||||||
|
|
||||||
// (a) Miss: the fabricated default is returned but not cached.
|
// (a) Miss: the fabricated default is returned but not cached.
|
||||||
let (bm, _) = sys
|
let (bm, _) = sys
|
||||||
@@ -1078,6 +1183,9 @@ mod tests {
|
|||||||
let mut persisted = BucketMetadata::new("absent-bucket");
|
let mut persisted = BucketMetadata::new("absent-bucket");
|
||||||
persisted.policy_config_json = b"persisted-marker".to_vec();
|
persisted.policy_config_json = b"persisted-marker".to_vec();
|
||||||
sys.persist_and_set(persisted).await.expect("metadata should persist");
|
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();
|
sys.metadata_map.write().await.clear();
|
||||||
let _ = sys
|
let _ = sys
|
||||||
.get_config("absent-bucket")
|
.get_config("absent-bucket")
|
||||||
@@ -1089,14 +1197,47 @@ mod tests {
|
|||||||
.expect("lazily loaded persisted metadata must be cached");
|
.expect("lazily loaded persisted metadata must be cached");
|
||||||
assert_eq!(cached.policy_config_json, b"persisted-marker".to_vec());
|
assert_eq!(cached.policy_config_json, b"persisted-marker".to_vec());
|
||||||
|
|
||||||
// (c) A refresh-load miss (no persisted metadata readable) must not
|
// (c) Persisted metadata left behind after physical deletion must not
|
||||||
// replace an existing entry.
|
// 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");
|
let mut kept = BucketMetadata::new("kept-bucket");
|
||||||
kept.policy_config_json = b"kept-marker".to_vec();
|
kept.policy_config_json = b"kept-marker".to_vec();
|
||||||
sys.set("kept-bucket".to_string(), Arc::new(kept)).await;
|
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 mut failed = HashSet::new();
|
||||||
let refresh_targets = vec!["kept-bucket".to_string()];
|
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
|
let kept = sys
|
||||||
.get("kept-bucket")
|
.get("kept-bucket")
|
||||||
.await
|
.await
|
||||||
@@ -1106,6 +1247,50 @@ mod tests {
|
|||||||
b"kept-marker".to_vec(),
|
b"kept-marker".to_vec(),
|
||||||
"a fabricated refresh default must not replace real metadata"
|
"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]
|
#[tokio::test]
|
||||||
@@ -1130,12 +1315,18 @@ mod tests {
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn update_config_with_persists_tagging_rewrite_across_disk_reload() {
|
async fn update_config_with_persists_tagging_rewrite_across_disk_reload() {
|
||||||
use crate::bucket::metadata::BUCKET_TAGGING_CONFIG;
|
use crate::bucket::metadata::BUCKET_TAGGING_CONFIG;
|
||||||
|
use crate::storage_api_contracts::bucket::MakeBucketOptions;
|
||||||
use s3s::dto::Tag;
|
use s3s::dto::Tag;
|
||||||
|
|
||||||
let (_dirs, ecstore) = isolated_store_over_temp_disks().await;
|
let (_dirs, ecstore) = isolated_store_over_temp_disks().await;
|
||||||
let mut sys = BucketMetadataSys::new(ecstore);
|
|
||||||
|
|
||||||
let bucket = "swift-tagging-bucket";
|
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))
|
sys.persist_and_set(BucketMetadata::new(bucket))
|
||||||
.await
|
.await
|
||||||
.expect("initial metadata should persist");
|
.expect("initial metadata should persist");
|
||||||
|
|||||||
@@ -89,6 +89,22 @@ fn reduce_pool_write_quorum_errs(per_pool_errs: &[Option<Error>]) -> Option<Erro
|
|||||||
reduce_write_quorum_errs(per_pool_errs, BUCKET_OP_IGNORED_ERRS, pool_write_quorum(per_pool_errs.len()))
|
reduce_write_quorum_errs(per_pool_errs, BUCKET_OP_IGNORED_ERRS, pool_write_quorum(per_pool_errs.len()))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn resolve_heal_bucket_mode(opts: &mut HealOpts, pool_errs: &[Option<Error>]) -> 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]
|
#[async_trait]
|
||||||
pub trait PeerS3Client: Debug + Sync + Send + 'static {
|
pub trait PeerS3Client: Debug + Sync + Send + 'static {
|
||||||
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem>;
|
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem>;
|
||||||
@@ -159,10 +175,7 @@ impl S3PeerSys {
|
|||||||
pool_errs.push(reduce_pool_write_quorum_errs(&per_pool_errs));
|
pool_errs.push(reduce_pool_write_quorum_errs(&per_pool_errs));
|
||||||
}
|
}
|
||||||
|
|
||||||
if !opts.recreate {
|
resolve_heal_bucket_mode(&mut opts, &pool_errs)?;
|
||||||
opts.remove = is_all_buckets_not_found(&pool_errs);
|
|
||||||
opts.recreate = !opts.remove;
|
|
||||||
}
|
|
||||||
|
|
||||||
let mut futures = Vec::new();
|
let mut futures = Vec::new();
|
||||||
let heal_bucket_results = Arc::new(RwLock::new(vec![HealResultItem::default(); self.clients.len()]));
|
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);
|
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]
|
#[tokio::test]
|
||||||
async fn test_make_bucket_reduces_quorum_by_pool_participants() {
|
async fn test_make_bucket_reduces_quorum_by_pool_participants() {
|
||||||
let peer_sys = S3PeerSys {
|
let peer_sys = S3PeerSys {
|
||||||
|
|||||||
@@ -87,7 +87,7 @@ fn bucket_deleted_marker_volume(bucket: &str) -> String {
|
|||||||
format!("{RUSTFS_META_BUCKET}/{}", bucket_deleted_marker_prefix(bucket))
|
format!("{RUSTFS_META_BUCKET}/{}", bucket_deleted_marker_prefix(bucket))
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn await_bucket_namespace_operation<T, F>(
|
pub(crate) async fn await_bucket_namespace_operation<T, F>(
|
||||||
guard: Option<&rustfs_lock::NamespaceLockGuard>,
|
guard: Option<&rustfs_lock::NamespaceLockGuard>,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
operation: &'static str,
|
operation: &'static str,
|
||||||
|
|||||||
@@ -141,6 +141,7 @@ fn should_enqueue_transition_immediately(oi: &ObjectInfo) -> bool {
|
|||||||
const MAX_UPLOADS_LIST: usize = 10000;
|
const MAX_UPLOADS_LIST: usize = 10000;
|
||||||
|
|
||||||
mod bucket;
|
mod bucket;
|
||||||
|
pub(crate) use bucket::await_bucket_namespace_operation;
|
||||||
mod heal;
|
mod heal;
|
||||||
mod heal_walk;
|
mod heal_walk;
|
||||||
pub use heal_walk::HealWalkVersion;
|
pub use heal_walk::HealWalkVersion;
|
||||||
|
|||||||
Reference in New Issue
Block a user