mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 12:31:36 +00:00
fix(ecstore): admit cold bucket metadata at read quorum (#8120)
* fix(ecstore): admit cold bucket metadata at read quorum Encode a durable creation-commit state in bucket metadata so committed Object Lock buckets can use read quorum for cold loads while pre-physical creation intents remain fail-closed. Drain commit-write fan-out, keep write quorum for uncommitted intents, and add regression coverage for exact, below, migrated, and partial-create quorum boundaries. * fix(ecstore): fence bucket creation-commit persistence Address review findings on the creation-commit proof. Persist the proof only under the bucket metadata transaction fence: at bucket creation, and through a fenced migration that re-reads the authoritative metadata and revalidates physical presence at write quorum before writing. This stops a stale snapshot from reverting an acknowledged configuration update or outliving a delete/recreate. Establish commitment when Object Lock is enabled on an existing bucket, inside the same configuration mutation, so a healthy cluster no longer rejects object operations with ErasureWriteQuorum. * fix(ecstore): keep commit fence error message stable The error(format!) ratchet requires a stable Display for quorum bucketing; use a fixed message instead of embedding the bucket name. * test(ecstore): reuse canonical Object Lock fixture in regression tests The s3s footprint ratchet is shrink-only. Use the existing ENABLED_OBJECT_LOCK_CONFIG static instead of naming s3s DTO types in store tests. --------- Co-authored-by: hector <42570491+majinghe@users.noreply.github.com> Co-authored-by: Hauser <housemecn@gmail.com>
This commit is contained in:
@@ -18,9 +18,10 @@ use super::versioning::VersioningApi;
|
||||
use super::{quota::BucketQuota, target::BucketTargets};
|
||||
use crate::bucket::replication::invalid_replication_config_status_field;
|
||||
use crate::bucket::utils::deserialize;
|
||||
use crate::config::com::{read_config, read_config_preserve_empty, save_config};
|
||||
use crate::config::com::{read_config, read_config_preserve_empty, save_config, save_config_with_opts};
|
||||
use crate::disk::BUCKET_META_PREFIX;
|
||||
use crate::error::{Error, Result};
|
||||
use crate::object_api::WriteCompletion;
|
||||
use crate::runtime::sources as runtime_sources;
|
||||
use crate::store::ECStore;
|
||||
use byteorder::{BigEndian, ByteOrder, LittleEndian};
|
||||
@@ -437,6 +438,7 @@ pub struct BucketMetadata {
|
||||
pub lock_enabled: bool, // While marked as unused, it may need to be retained
|
||||
pub bucket_incarnation_id: Uuid,
|
||||
pub(crate) bucket_incarnation_sidecar: bool,
|
||||
pub(crate) bucket_creation_committed: bool,
|
||||
pub policy_config_json: Vec<u8>,
|
||||
pub notification_config_xml: Vec<u8>,
|
||||
pub lifecycle_config_xml: Vec<u8>,
|
||||
@@ -511,6 +513,7 @@ impl Default for BucketMetadata {
|
||||
lock_enabled: Default::default(),
|
||||
bucket_incarnation_id: Uuid::nil(),
|
||||
bucket_incarnation_sidecar: false,
|
||||
bucket_creation_committed: false,
|
||||
policy_config_json: Default::default(),
|
||||
notification_config_xml: Default::default(),
|
||||
lifecycle_config_xml: Default::default(),
|
||||
@@ -594,6 +597,22 @@ impl BucketMetadata {
|
||||
metadata
|
||||
}
|
||||
|
||||
/// Whether this generation still needs its creation-commit proof.
|
||||
///
|
||||
/// The proof is only ever written while the caller holds the bucket
|
||||
/// metadata transaction fence; see
|
||||
/// `BucketMetadataSys::migrate_bucket_creation_commit`.
|
||||
///
|
||||
/// Metadata without an authoritative incarnation sidecar is owned by the
|
||||
/// legacy migration path, which mints the sidecar first; it must not be
|
||||
/// rewritten by this migration.
|
||||
pub(crate) fn needs_bucket_creation_commit(&self) -> bool {
|
||||
self.lock_enabled
|
||||
&& !self.bucket_creation_committed
|
||||
&& self.bucket_incarnation_sidecar
|
||||
&& !self.bucket_incarnation_id.is_nil()
|
||||
}
|
||||
|
||||
pub fn save_file_path(&self) -> String {
|
||||
format!("{}/{}/{}", BUCKET_META_PREFIX, self.name.as_str(), BUCKET_METADATA_FILE)
|
||||
}
|
||||
@@ -733,6 +752,7 @@ impl BucketMetadata {
|
||||
"TableBucketConfigUpdatedAt" => self.table_bucket_config_updated_at = read_msgp_time_value(rd)?,
|
||||
"DurabilityConfigUpdatedAt" => self.durability_config_updated_at = read_msgp_time_value(rd)?,
|
||||
"OnDemandMigrationConfigUpdatedAt" => self.on_demand_migration_config_updated_at = read_msgp_time_value(rd)?,
|
||||
"BucketCreationCommitted" => self.bucket_creation_committed = read_msgp_bool(rd)?,
|
||||
other => {
|
||||
tracing::debug!(field = %other, "BucketMetadata decode_from: skipping unknown field");
|
||||
skip_msgp_value(rd)?;
|
||||
@@ -745,8 +765,8 @@ impl BucketMetadata {
|
||||
|
||||
/// Encode to msgp bytes. Field order follows MinIO BucketMetadata for compatibility.
|
||||
pub fn encode_to<W: Write>(&self, wr: &mut W) -> Result<()> {
|
||||
// Map size: MinIO fields (25) + RustFS extensions (21)
|
||||
let map_len: u32 = 46;
|
||||
// Map size: MinIO fields (25) + RustFS extensions (22)
|
||||
let map_len: u32 = 47;
|
||||
rmp::encode::write_map_len(wr, map_len)?;
|
||||
|
||||
// MinIO field order (same as Go struct)
|
||||
@@ -827,6 +847,8 @@ impl BucketMetadata {
|
||||
write_msgp_time(wr, self.durability_config_updated_at)?;
|
||||
rmp::encode::write_str(wr, "OnDemandMigrationConfigUpdatedAt")?;
|
||||
write_msgp_time(wr, self.on_demand_migration_config_updated_at)?;
|
||||
rmp::encode::write_str(wr, "BucketCreationCommitted")?;
|
||||
rmp::encode::write_bool(wr, self.bucket_creation_committed)?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -982,6 +1004,11 @@ impl BucketMetadata {
|
||||
self.object_lock_config = None;
|
||||
if !data.is_empty() {
|
||||
self.lock_enabled = true;
|
||||
// Enabling Object Lock is only reachable for a bucket that
|
||||
// already exists: the caller validated physical presence and
|
||||
// holds the incarnation fence, so this generation is committed
|
||||
// and must never be mistaken for a pre-physical creation intent.
|
||||
self.bucket_creation_committed = true;
|
||||
}
|
||||
self.object_lock_config_xml = data;
|
||||
self.object_lock_config_updated_at = updated;
|
||||
@@ -1086,6 +1113,21 @@ impl BucketMetadata {
|
||||
/// server's bucket metadata lands in that server's `.rustfs.sys`, not the
|
||||
/// ambient (first) one. [`BucketMetadata::save`] keeps the ambient default.
|
||||
pub async fn save_with_store(&mut self, store: std::sync::Arc<crate::store::ECStore>) -> Result<()> {
|
||||
self.save_with_store_completion(store, WriteCompletion::Quorum).await
|
||||
}
|
||||
|
||||
/// Persist a committed creation record and drain the rename fan-out before
|
||||
/// returning. A quorum ACK is insufficient here: later degraded reads may
|
||||
/// have to reconstruct this record from any read quorum of surviving shards.
|
||||
pub(crate) async fn save_with_store_committed(&mut self, store: std::sync::Arc<crate::store::ECStore>) -> Result<()> {
|
||||
self.save_with_store_completion(store, WriteCompletion::TailDrained).await
|
||||
}
|
||||
|
||||
async fn save_with_store_completion(
|
||||
&mut self,
|
||||
store: std::sync::Arc<crate::store::ECStore>,
|
||||
write_completion: WriteCompletion,
|
||||
) -> Result<()> {
|
||||
self.parse_all_configs()?;
|
||||
let mut buf: Vec<u8> = vec![0; 4];
|
||||
|
||||
@@ -1099,7 +1141,17 @@ impl BucketMetadata {
|
||||
|
||||
buf.extend_from_slice(&data);
|
||||
|
||||
save_config(store, self.save_file_path().as_str(), buf).await?;
|
||||
save_config_with_opts(
|
||||
store,
|
||||
self.save_file_path().as_str(),
|
||||
buf,
|
||||
&crate::object_api::ObjectOptions {
|
||||
max_parity: true,
|
||||
write_completion,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -1451,8 +1503,13 @@ pub(crate) async fn load_bucket_metadata_parse_with_presence(
|
||||
}
|
||||
};
|
||||
|
||||
let incarnation = load_bucket_incarnation(api, bucket).await?;
|
||||
let incarnation = load_bucket_incarnation(api.clone(), bucket).await?;
|
||||
if persisted {
|
||||
if bm.name.is_empty() {
|
||||
bm.name = bucket.to_string();
|
||||
} else if bm.name != bucket {
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
if let Some(incarnation) = incarnation {
|
||||
if !bm.bucket_incarnation_id.is_nil() && bm.bucket_incarnation_id != incarnation {
|
||||
return Err(Error::other("bucket incarnation sidecar does not match bucket metadata"));
|
||||
@@ -1622,6 +1679,7 @@ mod test {
|
||||
BucketMetadata::check_header(&body).expect("recovered body is a valid .metadata.bin");
|
||||
let mut bm = BucketMetadata::unmarshal(&body[4..]).expect("unmarshal recovered blob");
|
||||
assert_eq!(bm.name, "interop");
|
||||
assert!(!bm.bucket_creation_committed, "legacy metadata must default to uncommitted");
|
||||
bm.parse_all_configs().expect("parse recovered configs");
|
||||
assert!(bm.lifecycle_config.is_some());
|
||||
}
|
||||
@@ -1630,13 +1688,15 @@ mod test {
|
||||
async fn marshal_msg() {
|
||||
// write_time(OffsetDateTime::UNIX_EPOCH).unwrap();
|
||||
|
||||
let bm = BucketMetadata::new("dada");
|
||||
let mut bm = BucketMetadata::new("dada");
|
||||
bm.bucket_creation_committed = true;
|
||||
|
||||
let buf = bm.marshal_msg().unwrap();
|
||||
|
||||
let new = BucketMetadata::unmarshal(&buf).unwrap();
|
||||
|
||||
assert_eq!(bm.name, new.name);
|
||||
assert!(new.bucket_creation_committed);
|
||||
assert!(!bm.bucket_incarnation_id.is_nil());
|
||||
assert_eq!(bm.bucket_incarnation_id, new.bucket_incarnation_id);
|
||||
}
|
||||
|
||||
@@ -527,11 +527,19 @@ pub(crate) async fn set_new_bucket_metadata_in(
|
||||
lock.persist_new_and_set(bm).await
|
||||
}
|
||||
|
||||
pub(crate) async fn cache_bucket_metadata_in(ctx: &crate::runtime::instance::InstanceContext, bm: BucketMetadata) -> Result<()> {
|
||||
pub(crate) async fn set_new_bucket_metadata_intent_in(
|
||||
ctx: &crate::runtime::instance::InstanceContext,
|
||||
bm: BucketMetadata,
|
||||
) -> Result<()> {
|
||||
let sys = bucket_metadata_sys_of(ctx)?;
|
||||
let lock = sys.read().await;
|
||||
lock.set(bm.name.clone(), Arc::new(bm)).await;
|
||||
Ok(())
|
||||
lock.persist_new_bucket_metadata_intent(bm).await
|
||||
}
|
||||
|
||||
pub(crate) async fn commit_bucket_metadata_in(ctx: &crate::runtime::instance::InstanceContext, bm: BucketMetadata) -> Result<()> {
|
||||
let sys = bucket_metadata_sys_of(ctx)?;
|
||||
let lock = sys.read().await;
|
||||
lock.commit_bucket_metadata(bm).await
|
||||
}
|
||||
|
||||
pub(crate) async fn remove_bucket_metadata_in(ctx: &crate::runtime::instance::InstanceContext, bucket: &str) -> Result<bool> {
|
||||
@@ -2074,9 +2082,32 @@ impl BucketMetadataSys {
|
||||
}
|
||||
|
||||
async fn persist_new_and_set(&self, mut bm: BucketMetadata) -> Result<()> {
|
||||
bm.bucket_creation_committed = true;
|
||||
bm.save_with_store_committed(self.object_store()).await?;
|
||||
save_bucket_incarnation(self.object_store(), &bm.name, bm.bucket_incarnation_id).await?;
|
||||
bm.bucket_incarnation_sidecar = true;
|
||||
self.set(bm.name.clone(), Arc::new(bm)).await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn persist_new_bucket_metadata_intent(&self, mut bm: BucketMetadata) -> Result<()> {
|
||||
bm.bucket_creation_committed = false;
|
||||
bm.save_with_store(self.object_store()).await?;
|
||||
save_bucket_incarnation(self.object_store(), &bm.name, bm.bucket_incarnation_id).await?;
|
||||
bm.bucket_incarnation_sidecar = true;
|
||||
// A pending creation intent must not enter the authoritative cache. The
|
||||
// caller publishes it only after physical creation and commit-marker write.
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn commit_bucket_metadata(&self, mut bm: BucketMetadata) -> Result<()> {
|
||||
bm.parse_all_configs()?;
|
||||
bm.bucket_incarnation_sidecar = true;
|
||||
// The caller holds the full bucket creation fence (lifecycle, metadata
|
||||
// transaction, and namespace write) and physical creation already
|
||||
// returned, so this generation is committed.
|
||||
bm.bucket_creation_committed = true;
|
||||
bm.save_with_store_committed(self.object_store()).await?;
|
||||
self.set(bm.name.clone(), Arc::new(bm)).await;
|
||||
Ok(())
|
||||
}
|
||||
@@ -2192,18 +2223,25 @@ impl BucketMetadataSys {
|
||||
)
|
||||
.await?;
|
||||
|
||||
let bm = Arc::new(bm);
|
||||
|
||||
if persisted {
|
||||
// Object Lock metadata is persisted before physical bucket creation as a
|
||||
// creation intent. The commit marker distinguishes a committed bucket
|
||||
// from an intent left behind by an interrupted create.
|
||||
let require_write_quorum = bm.needs_bucket_creation_commit();
|
||||
|
||||
await_bucket_namespace_operation(
|
||||
Some(&guard),
|
||||
bucket,
|
||||
"lazy bucket metadata existence check",
|
||||
Box::pin(async {
|
||||
self.object_store()
|
||||
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
|
||||
.await
|
||||
.map(|_| ())
|
||||
let store = self.object_store();
|
||||
let options = crate::storage_api_contracts::bucket::BucketOptions::default();
|
||||
let info = if require_write_quorum {
|
||||
store.get_bucket_info_from_sets(bucket, &options).await
|
||||
} else {
|
||||
store.get_bucket_info_from_sets_at_read_quorum(bucket, &options).await
|
||||
};
|
||||
info.map(|_| ())
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
@@ -2212,6 +2250,7 @@ impl BucketMetadataSys {
|
||||
"bucket namespace lock was lost before lazy bucket metadata publish: {bucket}"
|
||||
)));
|
||||
}
|
||||
let bm = Arc::new(bm);
|
||||
let _publish_guard = self
|
||||
.lock_metadata_publish(bucket, &guard, "lazy bucket metadata publish")
|
||||
.await?;
|
||||
@@ -2226,6 +2265,7 @@ impl BucketMetadataSys {
|
||||
sync_bucket_target_sys(bucket, &bm).await;
|
||||
sync_bucket_durability(bucket, &bm);
|
||||
sync_on_demand_migration(bucket, &bm);
|
||||
return Ok((bm, true));
|
||||
} else {
|
||||
let exists = self
|
||||
.bucket_exists(bucket, &guard, "lazy bucket metadata existence check")
|
||||
@@ -2245,7 +2285,7 @@ impl BucketMetadataSys {
|
||||
}
|
||||
}
|
||||
|
||||
Ok((bm, true))
|
||||
Ok((Arc::new(bm), true))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2430,7 +2470,13 @@ impl BucketMetadataSys {
|
||||
}
|
||||
|
||||
async fn get_metadata_authority(&self, bucket: &str) -> Result<BucketMetadataAuthority> {
|
||||
if let Some(bm) = self.metadata_map.read().await.get(bucket).cloned() {
|
||||
// Bind before matching: holding the map guard across the migration would
|
||||
// deadlock on the publish write.
|
||||
let cached = self.metadata_map.read().await.get(bucket).cloned();
|
||||
if let Some(bm) = cached {
|
||||
if bm.needs_bucket_creation_commit() {
|
||||
return self.migrate_bucket_creation_commit(bucket).await;
|
||||
}
|
||||
return Ok(BucketMetadataAuthority::Authoritative(bm));
|
||||
}
|
||||
if self.fabricated_metadata.read().await.contains(bucket) {
|
||||
@@ -2442,8 +2488,13 @@ impl BucketMetadataSys {
|
||||
|
||||
self.get_config(bucket).await?;
|
||||
|
||||
if let Some(bm) = self.metadata_map.read().await.get(bucket).cloned() {
|
||||
Ok(BucketMetadataAuthority::Authoritative(bm))
|
||||
let loaded = self.metadata_map.read().await.get(bucket).cloned();
|
||||
if let Some(bm) = loaded {
|
||||
if bm.needs_bucket_creation_commit() {
|
||||
self.migrate_bucket_creation_commit(bucket).await
|
||||
} else {
|
||||
Ok(BucketMetadataAuthority::Authoritative(bm))
|
||||
}
|
||||
} else if self.fabricated_metadata.read().await.contains(bucket) {
|
||||
Ok(BucketMetadataAuthority::Fabricated)
|
||||
} else if self.missing_buckets.get(bucket).await.is_some() {
|
||||
@@ -2453,6 +2504,70 @@ impl BucketMetadataSys {
|
||||
}
|
||||
}
|
||||
|
||||
/// Persist the creation-commit proof for an Object Lock bucket created
|
||||
/// before this field existed.
|
||||
///
|
||||
/// Lock order matches the established config-write order: metadata
|
||||
/// transaction (write), then bucket namespace (read). The authoritative
|
||||
/// metadata is re-read under both fences and the physical namespace is
|
||||
/// revalidated at write quorum, so a stale snapshot can neither revert an
|
||||
/// acknowledged configuration update nor outlive a delete/recreate.
|
||||
async fn migrate_bucket_creation_commit(&self, bucket: &str) -> Result<BucketMetadataAuthority> {
|
||||
let transaction_lock = self
|
||||
.object_store()
|
||||
.new_ns_lock(RUSTFS_META_BUCKET, &bucket_metadata_transaction_lock_key(bucket))
|
||||
.await?;
|
||||
let transaction_guard = transaction_lock
|
||||
.get_write_lock(crate::set_disk::get_lock_acquire_timeout())
|
||||
.await?;
|
||||
let namespace_lock = self.object_store().new_ns_lock(bucket, bucket).await?;
|
||||
let namespace_guard = namespace_lock
|
||||
.get_read_lock(crate::set_disk::get_lock_acquire_timeout())
|
||||
.await?;
|
||||
|
||||
// Write quorum must still prove that the physical namespace committed.
|
||||
await_bucket_namespace_operation(
|
||||
Some(&namespace_guard),
|
||||
bucket,
|
||||
"bucket creation commitment existence check",
|
||||
self.object_store()
|
||||
.get_bucket_info_from_sets(bucket, &crate::storage_api_contracts::bucket::BucketOptions::default()),
|
||||
)
|
||||
.await?;
|
||||
|
||||
let (mut metadata, persisted) = await_bucket_namespace_operation(
|
||||
Some(&namespace_guard),
|
||||
bucket,
|
||||
"bucket creation commitment metadata reload",
|
||||
load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true),
|
||||
)
|
||||
.await?;
|
||||
if !persisted {
|
||||
return Ok(BucketMetadataAuthority::Fabricated);
|
||||
}
|
||||
if metadata.needs_bucket_creation_commit() {
|
||||
metadata.bucket_creation_committed = true;
|
||||
metadata.save_with_store_committed(self.object_store()).await?;
|
||||
}
|
||||
|
||||
if transaction_guard.is_lock_lost() || namespace_guard.is_lock_lost() {
|
||||
// Stable message: formatted per-bucket detail would split quorum
|
||||
// error buckets (backlog#1845).
|
||||
return Err(Error::other("bucket creation commitment fence was lost before publish"));
|
||||
}
|
||||
let metadata = Arc::new(metadata);
|
||||
let _publish_guard = self
|
||||
.lock_metadata_publish(bucket, &namespace_guard, "bucket creation commitment publish")
|
||||
.await?;
|
||||
self.metadata_map
|
||||
.write()
|
||||
.await
|
||||
.insert(bucket.to_string(), Arc::clone(&metadata));
|
||||
self.fabricated_metadata.write().await.remove(bucket);
|
||||
self.missing_buckets.invalidate(bucket).await;
|
||||
Ok(BucketMetadataAuthority::Authoritative(metadata))
|
||||
}
|
||||
|
||||
async fn migrate_legacy_metadata(&self, bucket: &str) -> Result<BucketMetadataAuthority> {
|
||||
let transaction_lock = self
|
||||
.object_store()
|
||||
@@ -2564,7 +2679,6 @@ impl BucketMetadataSys {
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
|
||||
if namespace_guard.is_lock_lost() {
|
||||
return Err(Error::other(format!(
|
||||
"bucket namespace lock was lost before legacy metadata publish: {bucket}"
|
||||
@@ -2626,6 +2740,9 @@ impl BucketMetadataSys {
|
||||
load_bucket_metadata_parse_with_presence(self.object_store(), bucket, true),
|
||||
)
|
||||
.await?;
|
||||
if persisted && metadata.lock_enabled && !metadata.bucket_creation_committed {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
if persisted {
|
||||
Ok(BucketMetadataAuthority::Authoritative(Arc::new(metadata)))
|
||||
} else {
|
||||
|
||||
@@ -694,14 +694,13 @@ impl ECStore {
|
||||
let lock_newly_enabled = opts.lock_enabled && !meta.lock_enabled;
|
||||
if opts.lock_enabled {
|
||||
meta.lock_enabled = true;
|
||||
meta.object_lock_config_xml =
|
||||
crate::bucket::utils::serialize::<ObjectLockConfiguration>(&ENABLED_OBJECT_LOCK_CONFIG)?;
|
||||
meta.object_lock_config_xml = crate::bucket::utils::serialize(&*crate::store::ENABLED_OBJECT_LOCK_CONFIG)?;
|
||||
meta.versioning_config_xml = crate::bucket::utils::serialize::<VersioningConfiguration>(&ENABLED_VERSIONING_CONFIG)?;
|
||||
}
|
||||
|
||||
let metadata_persisted_before_physical = confirmed_missing && !is_meta_bucketname(bucket) && opts.lock_enabled;
|
||||
if metadata_persisted_before_physical {
|
||||
metadata_sys::set_new_bucket_metadata_in(&self.ctx, meta.clone()).await?;
|
||||
metadata_sys::set_new_bucket_metadata_intent_in(&self.ctx, meta.clone()).await?;
|
||||
if bucket_lifecycle_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|
||||
|| metadata_transaction_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|
||||
|| ns_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|
||||
@@ -759,12 +758,12 @@ impl ECStore {
|
||||
|
||||
let metadata_result = async {
|
||||
if metadata_persisted_before_physical {
|
||||
return Ok(());
|
||||
return metadata_sys::commit_bucket_metadata_in(&self.ctx, meta).await;
|
||||
}
|
||||
if is_meta_bucketname(bucket) {
|
||||
metadata_sys::set_bucket_metadata_in(&self.ctx, meta).await
|
||||
} else if existing_incarnation_is_authoritative && !lock_newly_enabled {
|
||||
metadata_sys::cache_bucket_metadata_in(&self.ctx, meta).await
|
||||
metadata_sys::commit_bucket_metadata_in(&self.ctx, meta).await
|
||||
} else {
|
||||
metadata_sys::set_new_bucket_metadata_in(&self.ctx, meta).await
|
||||
}
|
||||
@@ -910,6 +909,9 @@ impl ECStore {
|
||||
// permission to serve degraded reads.
|
||||
return Err(Error::ErasureReadQuorum);
|
||||
}
|
||||
if metadata.lock_enabled && !metadata.bucket_creation_committed {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
if metadata.name != bucket {
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
@@ -2393,6 +2395,342 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn lazy_metadata_load_uses_read_quorum_for_non_lock_bucket() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("lazy-read-quorum-{}", Uuid::new_v4().simple());
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("healthy namespace should accept bucket creation");
|
||||
|
||||
// Replace the instance metadata system with a cold cache while keeping the
|
||||
// persisted metadata and physical bucket in place.
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let set = &store.pools[0].disk_set[0];
|
||||
let offline = take_set_disks_offline(&store, set, &[0, 1]).await;
|
||||
|
||||
metadata_sys::get_on_demand_migration_config_in(&store, &bucket)
|
||||
.await
|
||||
.expect("cold metadata load should admit the exact namespace read quorum");
|
||||
assert!(
|
||||
metadata_sys::get_in(&store.ctx, &bucket).await.is_ok(),
|
||||
"the lazy load should publish authoritative metadata"
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.get_bucket_info_from_sets(&bucket, &BucketOptions::default())
|
||||
.await
|
||||
.expect_err("bucket mutations must retain their majority namespace check"),
|
||||
StorageError::ErasureWriteQuorum
|
||||
);
|
||||
|
||||
restore_set_disks(&store, set, offline).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn lazy_metadata_load_rejects_below_read_quorum() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("lazy-below-quorum-{}", Uuid::new_v4().simple());
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("healthy namespace should accept bucket creation");
|
||||
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let set = &store.pools[0].disk_set[0];
|
||||
let offline = take_set_disks_offline(&store, set, &[0, 1, 2]).await;
|
||||
|
||||
let error = metadata_sys::get_on_demand_migration_config_in(&store, &bucket)
|
||||
.await
|
||||
.expect_err("cold metadata load must reject read quorum minus one");
|
||||
assert!(
|
||||
matches!(error, StorageError::ErasureReadQuorum | StorageError::InsufficientReadQuorum(_, _)),
|
||||
"below read quorum must fail with a read-quorum error, got {error}"
|
||||
);
|
||||
assert!(
|
||||
metadata_sys::get_in(&store.ctx, &bucket).await.is_err(),
|
||||
"a below-quorum load must not publish metadata"
|
||||
);
|
||||
|
||||
restore_set_disks(&store, set, offline).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn lazy_metadata_load_keeps_lock_creation_intent_fail_closed() {
|
||||
let (temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("lazy-lock-intent-{}", Uuid::new_v4().simple());
|
||||
let mut intent = BucketMetadata::new(&bucket);
|
||||
intent.lock_enabled = true;
|
||||
metadata_sys::set_new_bucket_metadata_intent_in(&store.ctx, intent)
|
||||
.await
|
||||
.expect("Object Lock creation intent should persist metadata and its incarnation sidecar");
|
||||
|
||||
let set = &store.pools[0].disk_set[0];
|
||||
let offline = take_set_disks_offline(&store, set, &[0, 1]).await;
|
||||
assert_eq!(
|
||||
store
|
||||
.make_bucket_on_sets(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect_err("physical creation must not commit below write quorum"),
|
||||
StorageError::ErasureWriteQuorum
|
||||
);
|
||||
assert!(!temp_dir.path().join("pool0-disk0").join(&bucket).exists());
|
||||
assert!(!temp_dir.path().join("pool0-disk1").join(&bucket).exists());
|
||||
assert!(temp_dir.path().join("pool0-disk2").join(&bucket).is_dir());
|
||||
assert!(temp_dir.path().join("pool0-disk3").join(&bucket).is_dir());
|
||||
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let error = metadata_sys::get_on_demand_migration_config_in(&store, &bucket)
|
||||
.await
|
||||
.expect_err("Object Lock intent must not become visible at read quorum");
|
||||
assert_eq!(error, StorageError::ErasureWriteQuorum);
|
||||
assert!(
|
||||
metadata_sys::get_in(&store.ctx, &bucket).await.is_err(),
|
||||
"the rejected Object Lock intent must not enter the metadata cache"
|
||||
);
|
||||
assert_eq!(
|
||||
store
|
||||
.get_bucket_info(&bucket, &BucketOptions::default())
|
||||
.await
|
||||
.expect_err("bucket-info read fallback must also reject a pending Object Lock intent"),
|
||||
StorageError::ErasureWriteQuorum
|
||||
);
|
||||
|
||||
restore_set_disks(&store, set, offline).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn lazy_metadata_load_admits_committed_lock_bucket_at_read_quorum() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("lazy-committed-lock-{}", Uuid::new_v4().simple());
|
||||
store
|
||||
.make_bucket(
|
||||
&bucket,
|
||||
&MakeBucketOptions {
|
||||
lock_enabled: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("healthy namespace should accept Object Lock bucket creation");
|
||||
let (metadata, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("read committed Object Lock metadata");
|
||||
assert!(persisted && metadata.bucket_creation_committed);
|
||||
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let set = &store.pools[0].disk_set[0];
|
||||
let offline = take_set_disks_offline(&store, set, &[0, 1]).await;
|
||||
|
||||
metadata_sys::get_on_demand_migration_config_in(&store, &bucket)
|
||||
.await
|
||||
.expect("committed Object Lock metadata should load at read quorum");
|
||||
assert!(
|
||||
metadata_sys::get_in(&store.ctx, &bucket).await.is_ok(),
|
||||
"committed Object Lock metadata should be published"
|
||||
);
|
||||
|
||||
restore_set_disks(&store, set, offline).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn creation_commit_backfill_runs_under_the_config_write_fence() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("lazy-commit-backfill-{}", Uuid::new_v4().simple());
|
||||
let mut intent = BucketMetadata::new(&bucket);
|
||||
intent.lock_enabled = true;
|
||||
metadata_sys::set_new_bucket_metadata_intent_in(&store.ctx, intent)
|
||||
.await
|
||||
.expect("persist pending Object Lock intent");
|
||||
store
|
||||
.make_bucket_on_sets(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("physical creation should commit at full write quorum");
|
||||
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
metadata_sys::get_object_lock_config_state_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("the fenced authority path should backfill the creation commit");
|
||||
let (metadata, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("read backfilled creation commit");
|
||||
assert!(persisted && metadata.bucket_creation_committed);
|
||||
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let set = &store.pools[0].disk_set[0];
|
||||
let offline = take_set_disks_offline(&store, set, &[0, 1]).await;
|
||||
metadata_sys::get_on_demand_migration_config_in(&store, &bucket)
|
||||
.await
|
||||
.expect("backfilled Object Lock bucket should load at read quorum");
|
||||
restore_set_disks(&store, set, offline).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn enabling_object_lock_on_existing_bucket_establishes_commitment() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("enable-lock-legacy-{}", Uuid::new_v4().simple());
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create a versioned non-lock bucket");
|
||||
|
||||
// Simulate a bucket persisted by a build before the commit field existed.
|
||||
let mut legacy = metadata_sys::get_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("cached metadata")
|
||||
.as_ref()
|
||||
.clone();
|
||||
legacy.lock_enabled = false;
|
||||
legacy.object_lock_config_xml = Vec::new();
|
||||
legacy.object_lock_config = None;
|
||||
legacy.bucket_creation_committed = false;
|
||||
metadata_sys::set_new_bucket_metadata_intent_in(&store.ctx, legacy)
|
||||
.await
|
||||
.expect("persist legacy-shaped metadata");
|
||||
|
||||
let lock_xml = crate::bucket::utils::serialize(&*crate::store::ENABLED_OBJECT_LOCK_CONFIG)
|
||||
.expect("serialize Object Lock configuration");
|
||||
metadata_sys::update_in(&store.ctx, &bucket, "object-lock.xml", lock_xml)
|
||||
.await
|
||||
.expect("enabling Object Lock on an existing bucket must succeed");
|
||||
|
||||
assert!(
|
||||
matches!(
|
||||
metadata_sys::get_object_lock_config_state_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("read the freshly enabled Object Lock state"),
|
||||
metadata_sys::ObjectLockConfigState::Configured { .. }
|
||||
),
|
||||
"a healthy cluster must not report ErasureWriteQuorum after enabling Object Lock"
|
||||
);
|
||||
let (metadata, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("read persisted metadata");
|
||||
assert!(persisted && metadata.bucket_creation_committed);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn creation_commit_backfill_preserves_concurrent_config_update() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("commit-backfill-pack-{}", Uuid::new_v4().simple());
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create bucket");
|
||||
|
||||
// The cached snapshot on the migrating node still predates another
|
||||
// node's acknowledged configuration update.
|
||||
let mut stale = metadata_sys::get_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("cached metadata")
|
||||
.as_ref()
|
||||
.clone();
|
||||
stale.bucket_creation_committed = false;
|
||||
metadata_sys::set_new_bucket_metadata_intent_in(&store.ctx, stale.clone())
|
||||
.await
|
||||
.expect("persist the stale pre-upgrade snapshot");
|
||||
|
||||
// A newer, still-uncommitted copy carrying an acknowledged Object Lock
|
||||
// update lands on disk after the migrating node read its snapshot.
|
||||
let mut newer = stale.clone();
|
||||
newer.lock_enabled = true;
|
||||
newer.object_lock_config_xml = crate::bucket::utils::serialize(&*crate::store::ENABLED_OBJECT_LOCK_CONFIG)
|
||||
.expect("serialize Object Lock configuration");
|
||||
newer
|
||||
.save_with_store(store.clone())
|
||||
.await
|
||||
.expect("persist the newer configuration");
|
||||
|
||||
// A cold node must re-read the newest authoritative metadata under the
|
||||
// fence instead of rewriting its stale snapshot.
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
metadata_sys::get_object_lock_config_state_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("fenced commit migration");
|
||||
|
||||
let (metadata, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("read persisted metadata");
|
||||
assert!(persisted && metadata.bucket_creation_committed);
|
||||
assert!(
|
||||
metadata.lock_enabled && !metadata.object_lock_config_xml.is_empty(),
|
||||
"the acknowledged Object Lock update must survive the backfill"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn creation_commit_backfill_keeps_recreated_generation() {
|
||||
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], Some(2)).await;
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
let bucket = format!("commit-backfill-recreate-{}", Uuid::new_v4().simple());
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create bucket");
|
||||
let old_incarnation = store.bucket_incarnation_id_from_disk(&bucket).await.expect("old incarnation");
|
||||
|
||||
// A node that still holds the pre-recreation snapshot must never let it
|
||||
// outlive the replacement generation.
|
||||
let mut stale = metadata_sys::get_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("cached metadata")
|
||||
.as_ref()
|
||||
.clone();
|
||||
stale.lock_enabled = true;
|
||||
stale.object_lock_config_xml = crate::bucket::utils::serialize(&*crate::store::ENABLED_OBJECT_LOCK_CONFIG)
|
||||
.expect("serialize Object Lock configuration");
|
||||
stale.bucket_creation_committed = false;
|
||||
stale.save_with_store(store.clone()).await.expect("persist stale metadata");
|
||||
|
||||
store
|
||||
.delete_bucket(&bucket, &crate::storage_api_contracts::bucket::DeleteBucketOptions::default())
|
||||
.await
|
||||
.expect("delete bucket");
|
||||
store
|
||||
.make_bucket(
|
||||
&bucket,
|
||||
&MakeBucketOptions {
|
||||
lock_enabled: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("recreate bucket with Object Lock");
|
||||
let new_incarnation = store.bucket_incarnation_id_from_disk(&bucket).await.expect("new incarnation");
|
||||
assert_ne!(old_incarnation, new_incarnation);
|
||||
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
metadata_sys::get_object_lock_config_state_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("fenced backfill on the replacement generation");
|
||||
|
||||
let (metadata, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&store.ctx, &bucket)
|
||||
.await
|
||||
.expect("read persisted metadata");
|
||||
assert!(persisted);
|
||||
assert_eq!(
|
||||
metadata.bucket_incarnation_id, new_incarnation,
|
||||
"the replacement generation must stay authoritative"
|
||||
);
|
||||
assert!(metadata.bucket_creation_committed);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn bucket_info_read_quorum_is_scoped_to_each_erasure_set() {
|
||||
|
||||
Reference in New Issue
Block a user