From 212b20f89097369176166a6437fce38f47e1b938 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 26 Sep 2026 21:40:40 +0800 Subject: [PATCH] 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 --- crates/ecstore/src/bucket/metadata.rs | 72 ++++- crates/ecstore/src/bucket/metadata_sys.rs | 145 ++++++++- crates/ecstore/src/store/bucket.rs | 348 +++++++++++++++++++++- 3 files changed, 540 insertions(+), 25 deletions(-) diff --git a/crates/ecstore/src/bucket/metadata.rs b/crates/ecstore/src/bucket/metadata.rs index 2407a7e5b..cc75e2d7f 100644 --- a/crates/ecstore/src/bucket/metadata.rs +++ b/crates/ecstore/src/bucket/metadata.rs @@ -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, pub notification_config_xml: Vec, pub lifecycle_config_xml: Vec, @@ -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(&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) -> 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) -> Result<()> { + self.save_with_store_completion(store, WriteCompletion::TailDrained).await + } + + async fn save_with_store_completion( + &mut self, + store: std::sync::Arc, + write_completion: WriteCompletion, + ) -> Result<()> { self.parse_all_configs()?; let mut buf: Vec = 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); } diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 255715ff8..a509c87d4 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -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 { @@ -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 { - 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 { + 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 { 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 { diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index 72553dc43..fdd539f4a 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -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::(&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::(&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() {