From 00aeb1291470bb6a0f2099cb62e6d1a91b6d2d26 Mon Sep 17 00:00:00 2001 From: cxymds Date: Thu, 10 Sep 2026 20:14:07 +0800 Subject: [PATCH] fix(ilm): preserve cleanup ownership on tiered overwrites (#7639) * fix(ilm): preserve cleanup ownership on tiered overwrites * fix(ci): refresh E2E selection for tier overwrite regression --- .config/e2e-full-selection.txt | 4 +- .config/e2e-smoke-selection.txt | 2 +- crates/e2e_test/src/reliant/tiering.rs | 79 ++++- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 71 +++- .../src/set_disk/core/io_primitives.rs | 89 ++++- crates/ecstore/src/set_disk/ops/object.rs | 103 ++++++ crates/ecstore/src/store/init.rs | 218 ++++++++++++ crates/filemeta/src/filemeta.rs | 334 +++++++++++++++++- .../ilm-tiering-persistence-contracts.md | 12 + 9 files changed, 904 insertions(+), 8 deletions(-) diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index 1ab608350..5890ebce7 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5 -sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e +sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5 +sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index aa3804f21..f3b450db0 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0 +sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9 diff --git a/crates/e2e_test/src/reliant/tiering.rs b/crates/e2e_test/src/reliant/tiering.rs index 4f97f724e..f013bf2f2 100644 --- a/crates/e2e_test/src/reliant/tiering.rs +++ b/crates/e2e_test/src/reliant/tiering.rs @@ -52,8 +52,8 @@ use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{ BucketLifecycleConfiguration, BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, ExpirationStatus, - LifecycleRule, LifecycleRuleFilter, NoncurrentVersionTransition, RestoreRequest, Transition, TransitionStorageClass, - VersioningConfiguration, + LifecycleRule, LifecycleRuleFilter, MetadataDirective, NoncurrentVersionTransition, RestoreRequest, Transition, + TransitionStorageClass, VersioningConfiguration, }; use http::Method; use serde::Deserialize; @@ -963,6 +963,81 @@ async fn test_hermetic_transition_main_path() -> TestResult { Ok(()) } +/// PUT and materialized self-copy must retain cleanup ownership of a replaced +/// transitioned null version while publishing the new bytes and metadata. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn test_hermetic_transition_overwrite_and_self_copy() -> TestResult { + let mut cold = RustFSTestEnvironment::new().await?; + cold.access_key = "coldtieradmin".to_string(); + cold.secret_key = "coldtiersecret".to_string(); + cold.start_rustfs_server_without_cleanup(vec![]).await?; + let cold_client = cold.create_s3_client(); + cold_client.create_bucket().bucket(TIER_BUCKET).send().await?; + + let mut hot = RustFSTestEnvironment::new().await?; + start_tier_source(&mut hot, crate::common::FAST_DATA_USAGE_SCANNER_ENV).await?; + let hot_client = hot.create_s3_client(); + add_rustfs_tier(&hot, &cold).await?; + hot_client.create_bucket().bucket(SOURCE_BUCKET).send().await?; + + let data = payload(); + for self_copy in [false, true] { + hot_client + .put_bucket_lifecycle_configuration() + .bucket(SOURCE_BUCKET) + .lifecycle_configuration(BucketLifecycleConfiguration::builder().rules(transition_rule()?).build()?) + .send() + .await?; + put_multipart_object(&hot_client, SOURCE_BUCKET, OBJECT_KEY, &data).await?; + wait_for_transition(&hot_client, SOURCE_BUCKET, OBJECT_KEY, StdDuration::from_secs(90)).await?; + assert_eq!(cold_tier_object_count(&cold_client).await?, 1); + // Keep the replacement local so disappearance of the old remote + // object cannot be confused with another automatic transition. + hot_client.delete_bucket_lifecycle().bucket(SOURCE_BUCKET).send().await?; + + let expected = if self_copy { data.clone() } else { vec![0x73; 513] }; + if self_copy { + hot_client + .copy_object() + .bucket(SOURCE_BUCKET) + .key(OBJECT_KEY) + .copy_source(format!("{SOURCE_BUCKET}/{}", urlencoding::encode(OBJECT_KEY))) + .metadata_directive(MetadataDirective::Replace) + .content_type("text/plain") + .metadata("replacement", "kept") + .send() + .await?; + } else { + hot_client + .put_object() + .bucket(SOURCE_BUCKET) + .key(OBJECT_KEY) + .body(ByteStream::from(expected.clone())) + .content_type("text/plain") + .metadata("replacement", "kept") + .send() + .await?; + } + wait_for_cold_tier_empty(&cold_client, StdDuration::from_secs(90)).await?; + let current = hot_client.get_object().bucket(SOURCE_BUCKET).key(OBJECT_KEY).send().await?; + assert_eq!(current.content_type(), Some("text/plain")); + assert_eq!(current.metadata().and_then(|m| m.get("replacement")).map(String::as_str), Some("kept")); + assert!( + current + .metadata() + .is_none_or(|metadata| !metadata.contains_key(USER_META_KEY)) + ); + assert_eq!(current.body.collect().await?.into_bytes().as_ref(), expected.as_slice()); + hot_client + .delete_object() + .bucket(SOURCE_BUCKET) + .key(OBJECT_KEY) + .send() + .await?; + } + Ok(()) +} + /// Restore a transitioned object through a real RustFS remote tier. /// /// The test covers the externally visible copy-back contract that a mock tier diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 8b33231af..301c022a1 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -822,7 +822,7 @@ async fn scan_exact_free_version_targets( let mut targets = Vec::new(); for pool in &api.pools { for set in &pool.disk_set { - let versions = match set.load_file_info_versions_exact(&oi.bucket, &oi.name).await { + let versions = match set.load_file_info_versions_for_tier_cleanup(&oi.bucket, &oi.name).await { Ok(Some(versions)) => versions, Ok(None) => continue, Err(err) if is_err_strict_volume_not_found(&err) => continue, @@ -8100,6 +8100,75 @@ mod tests { } } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial] + async fn tier_overwrite_cleanup_retains_a_minority_live_remote_reference() { + let (disk_paths, ecstore) = setup_test_env().await; + let bucket = format!("overwrite-minority-{}", Uuid::new_v4()); + let object = "still-referenced"; + create_test_bucket(&ecstore, &bucket).await; + let (backend, identity) = register_recovery_mock_tier(&ecstore).await; + seed_recoverable_free_version(&disk_paths, &bucket, object, None, Some(identity.clone())).await; + let page = list_tier_free_versions(Arc::clone(&ecstore), 100, None, None, CancellationToken::new()) + .await + .expect("list persisted cleanup owner"); + let owner = page.items.into_iter().find(|oi| oi.bucket == bucket).expect("seeded owner"); + backend + .set_put_remote_version(Some(owner.transitioned_object.version_id.clone())) + .await; + let lease = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM") + .await + .expect("remote fixture lease"); + lease + .put( + &owner.transitioned_object.name, + rustfs_s3_client::transition_api::ReaderImpl::Body(bytes::Bytes::from_static(b"old")), + 3, + ) + .await + .expect("seed referenced remote bytes"); + drop(lease); + let path = disk_paths[0].join(&bucket).join(object).join(STORAGE_FORMAT_FILE); + let cleanup_metadata = fs::read(&path).await.expect("save completed replica"); + let mut live = FileInfo::new(object, 2, 2); + live.volume = bucket.clone(); + live.erasure.index = 1; + live.data_dir = Some(Uuid::new_v4()); + live.mod_time = Some(OffsetDateTime::now_utc()); + live.size = 3; + live.add_object_part(1, "149603e6c03516362a8da23f624db945".to_string(), 3, live.mod_time, 3, None, None); + live.transition_status = TRANSITION_COMPLETE.to_string(); + live.transition_tier = "WARM".to_string(); + live.transitioned_objname = owner.transitioned_object.name.clone(); + live.transition_version = Some(owner.transitioned_object.version_id.clone()); + live.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact; + rustfs_utils::http::insert_str(&mut live.metadata, rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity); + let mut old_metadata = FileMeta::new(); + old_metadata.add_version(live).expect("prepare minority live source"); + fs::write(&path, old_metadata.marshal_msg().expect("encode live source")) + .await + .expect("model one replica retained by an interrupted overwrite"); + + let err = super::cleanup_free_version_exact(Arc::clone(&ecstore), &owner, &CancellationToken::new()) + .await + .expect_err("quorum free versions cannot erase a minority live reference"); + assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); + assert_eq!(backend.remove_count().await, 0); + assert!(backend.contains(&owner.transitioned_object.name).await); + + fs::write(&path, cleanup_metadata) + .await + .expect("complete replica convergence"); + assert!( + super::cleanup_free_version_exact(Arc::clone(&ecstore), &owner, &CancellationToken::new()) + .await + .expect("converged cleanup can delete the exact remote owner") + ); + assert_eq!(backend.remove_count().await, 1); + assert!(!backend.contains(&owner.transitioned_object.name).await); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial] diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index e21969917..b74a397e4 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -3076,6 +3076,23 @@ impl SetDisks { &self, bucket: &str, object: &str, + ) -> Result> { + self.load_file_info_versions_for_cleanup(bucket, object, false).await + } + + pub(crate) async fn load_file_info_versions_for_tier_cleanup( + &self, + bucket: &str, + object: &str, + ) -> Result> { + self.load_file_info_versions_for_cleanup(bucket, object, true).await + } + + async fn load_file_info_versions_for_cleanup( + &self, + bucket: &str, + object: &str, + retain_unconfirmed_tier_references: bool, ) -> Result> { let disk_object = rustfs_utils::path::encode_dir_object(object); let disks = self.get_disks_internal().await; @@ -3152,12 +3169,24 @@ impl SetDisks { ))); } - let file_info_versions = FileMeta { + let mut file_info_versions = FileMeta { versions, ..Default::default() } .get_all_file_info_versions(bucket, object, true) .map_err(decode_error)?; + if retain_unconfirmed_tier_references { + // A failed overwrite may leave its live source on a minority + // of disks. Preserve that reference even if quorum merging + // selects only the replacement and its cleanup owner. + file_info_versions.versions.extend( + transition_copies + .into_values() + .flatten() + .map(|(version, _)| version) + .filter(|version| !version.tier_free_version()), + ); + } for file_info in file_info_versions .versions @@ -12214,6 +12243,64 @@ mod tests { assert!(result.is_err(), "missing disks must prevent metadata write quorum"); } + #[tokio::test] + async fn tier_overwrite_cleanup_rejects_unreadable_disk_despite_metadata_quorum() { + let bucket = "tier-unreadable-disk"; + let object = "object"; + let mut dirs = Vec::new(); + let mut disks = Vec::new(); + let mut fi = metadata_test_fileinfo(object); + fi.mod_time = Some(OffsetDateTime::now_utc()); + for index in 1..=3 { + let (dir, disk) = read_multiple_test_disk(bucket, &[]).await; + fi.erasure.index = index; + disk.write_metadata(bucket, bucket, object, fi.clone()) + .await + .expect("seed metadata quorum"); + dirs.push(dir); + disks.push(Some(disk)); + } + disks.push(None); + let set = io_primitives_test_set(disks, 2).await; + assert!( + set.load_file_info_versions_exact(bucket, object).await.is_err(), + "exact reads must preserve release's unreadable-replica fence" + ); + assert!( + set.load_file_info_versions_for_tier_cleanup(bucket, object).await.is_err(), + "unreadable replica may still reference the old remote object" + ); + } + + #[tokio::test] + async fn tier_overwrite_cleanup_rejects_minority_metadata_in_an_absent_set() { + let bucket = "tier-minority-metadata"; + let object = "object"; + let mut dirs = Vec::new(); + let mut disks = Vec::new(); + for index in 1..=4 { + let (dir, disk) = read_multiple_test_disk(bucket, &[]).await; + if index == 1 { + let mut fi = metadata_test_fileinfo(object); + fi.mod_time = Some(OffsetDateTime::now_utc()); + disk.write_metadata(bucket, bucket, object, fi) + .await + .expect("seed minority metadata"); + } + dirs.push(dir); + disks.push(Some(disk)); + } + let set = io_primitives_test_set(disks, 2).await; + assert!( + set.load_file_info_versions_exact(bucket, object).await.is_err(), + "exact reads must preserve release's minority-ownership fence" + ); + assert!( + set.load_file_info_versions_for_tier_cleanup(bucket, object).await.is_err(), + "absence on a majority cannot prove this physical set has no remote reference" + ); + } + #[tokio::test] async fn load_file_info_versions_exact_returns_versions_from_read_quorum() { let bucket = "exact-versions-bucket"; diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 9cfcd4e4d..47ac15613 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -3821,6 +3821,13 @@ impl SetDisks { } fi.metadata = user_defined; + if fi.version_id.is_none_or(|id| id.is_nil()) && !opts.data_movement && expected_restore_operation_id.is_none() { + // Every disk must publish the same cleanup owner alongside a + // replaced null version. This transient key is not persisted + // on the new object; recovery discovers the free-version in + // the committed xl.meta even if this request is cancelled. + fi.set_tier_free_version_id(&Uuid::new_v4().to_string()); + } fi.mod_time = mod_time; fi.size = w_size as i64; fi.versioned = opts.versioned || opts.version_suspended; @@ -18174,6 +18181,102 @@ mod put_object_tmp_cleanup_tests { drop(temp_dirs); } + #[tokio::test] + #[serial_test::serial(capacity_dirty_scope)] + async fn tier_overwrite_failed_quorum_and_cancellation_preserve_live_source() { + for cancel_before_rename in [false, true] { + let (dirs, disks, set) = hermetic_set_disks(4).await; + let bucket = "tier-overwrite-failure"; + let object = "still-live"; + make_completion_test_bucket(&disks, bucket).await; + let old_body = vec![0x31; TEST_OBJECT_SIZE]; + let mut metadata = HashMap::from([( + "x-amz-restore".to_string(), + "ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(), + )]); + for (suffix, value) in [ + (rustfs_utils::http::SUFFIX_TRANSITION_STATUS, "complete".to_string()), + (rustfs_utils::http::SUFFIX_TRANSITION_TIER, "WARM".to_string()), + (rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME, "remote/still-live".to_string()), + (rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, "exact-live-version".to_string()), + (rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE, "exact".to_string()), + (rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, "ab".repeat(32)), + ] { + rustfs_utils::http::insert_str(&mut metadata, suffix, value); + } + set.put_object( + bucket, + object, + &mut PutObjReader::from_vec(old_body.clone()), + &ObjectOptions { + user_defined: metadata, + write_completion: WriteCompletion::TailDrained, + ..Default::default() + }, + ) + .await + .expect("seed live transitioned source"); + wait_for_tmp_workspace_to_drain(&dirs, "seed write must drain").await; + let before = set + .load_file_info_versions_exact(bucket, object) + .await + .expect("read original metadata") + .expect("original exists"); + + if cancel_before_rename { + let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::AfterQuotaReservation); + let writer = Arc::clone(&set); + let put = tokio::spawn(async move { + writer + .put_object( + bucket, + object, + &mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]), + &ObjectOptions::default(), + ) + .await + }); + barrier.wait_until_paused().await; + put.abort(); + assert!(put.await.expect_err("cancel paused replacement").is_cancelled()); + wait_for_tmp_workspace_to_drain(&dirs, "cancelled replacement must roll back").await; + drop(barrier); + } else { + let _fault = rename_fault_injection::fail_rename_on(object, &[2, 3]); + let err = set + .put_object( + bucket, + object, + &mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]), + &ObjectOptions { + write_completion: WriteCompletion::TailDrained, + ..Default::default() + }, + ) + .await + .expect_err("two disk commits cannot satisfy write quorum three"); + assert!(matches!(err, Error::ErasureWriteQuorum | Error::InsufficientWriteQuorum(_, _)), "{err}"); + } + let after = set + .load_file_info_versions_exact(bucket, object) + .await + .expect("read rolled-back metadata") + .expect("live source must survive"); + assert_eq!(after.versions, before.versions, "failed replacement must preserve the live version"); + assert_eq!( + after.free_versions, before.free_versions, + "failed replacement must not publish a cleanup owner" + ); + let mut reader = set + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("live source remains readable"); + let mut actual = Vec::new(); + reader.stream.read_to_end(&mut actual).await.expect("read original bytes"); + assert_eq!(actual, old_body); + } + } + #[tokio::test] async fn cooperative_cancellation_while_waiting_for_namespace_lock_cleans_tmp_workspace() { let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index c2e087acc..a12b887ce 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -2832,6 +2832,224 @@ mod tests { shutdown.cancel(); } + #[cfg(feature = "test-util")] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[serial_test::serial(storage_class_env)] + async fn tier_overwrite_put_and_self_copy_recover_persisted_cleanup_owners() { + use crate::bucket::lifecycle::bucket_lifecycle_ops::ExpiryState; + use crate::bucket::lifecycle::tier_free_version_recovery::recover_tier_free_versions; + use rustfs_filemeta::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull}; + use rustfs_s3_client::transition_api::ReaderImpl; + use rustfs_utils::http::{ + SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID, + SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, insert_str, + }; + + let temp_dir = tempfile::tempdir().expect("create tier overwrite store"); + let (mut ctx, mut store, mut shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite", &[4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await; + let tier = "OVERWRITE-TIER"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier).await; + let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier) + .await + .expect("tier identity"); + let identity = rustfs_utils::crypto::hex(lease.backend_identity()); + drop(lease); + + for state in [Exact, KnownDisabled, SuspendedNull] { + for suspended in [false, true] { + for self_copy in [false, true] { + let bucket = format!("tier-overwrite-{}", Uuid::new_v4()); + let object = "object"; + let remote = format!("remote/{bucket}"); + let version = match state { + Exact => "opaque-overwrite-version", + SuspendedNull => "null", + _ => "", + }; + let payload = vec![0x5b; if suspended { 512 * 1024 } else { 257 }]; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create bucket"); + backend.set_put_remote_version(Some(version.to_string())).await; + let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier) + .await + .expect("seed tier lease"); + lease + .put( + &remote, + ReaderImpl::Body(bytes::Bytes::from(payload.clone())), + payload.len().try_into().expect("payload size"), + ) + .await + .expect("seed remote bytes"); + drop(lease); + let mut metadata = HashMap::from([ + ("content-type".to_string(), "application/octet-stream".to_string()), + ( + "x-amz-restore".to_string(), + "ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(), + ), + ]); + for (suffix, value) in [ + (SUFFIX_TRANSITION_STATUS, "complete"), + (SUFFIX_TRANSITION_TIER, tier), + (SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity.as_str()), + (SUFFIX_TRANSITIONED_OBJECTNAME, remote.as_str()), + (SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str()), + ] { + insert_str(&mut metadata, suffix, value.to_string()); + } + if !version.is_empty() { + insert_str(&mut metadata, SUFFIX_TRANSITIONED_VERSION_ID, version.to_string()); + } + let options = ObjectOptions { + version_suspended: suspended, + ..Default::default() + }; + store + .put_object( + &bucket, + object, + &mut PutObjReader::from_vec(payload.clone()), + &ObjectOptions { + user_defined: metadata, + ..options.clone() + }, + ) + .await + .expect("seed transitioned source with locally restored bytes"); + let expected = if self_copy { + payload.clone() + } else { + vec![0x73; payload.len()] + }; + let new_metadata = HashMap::from([ + ("content-type".to_string(), "text/plain".to_string()), + ("x-amz-meta-replacement".to_string(), "kept".to_string()), + ]); + if self_copy { + let mut source = store + .get_object_info(&bucket, object, &options) + .await + .expect("self-copy source"); + source.metadata_only = false; + source.user_defined = Arc::new(new_metadata); + source.put_object_reader = Some(PutObjReader::from_vec(expected.clone())); + store + .copy_object(&bucket, object, &bucket, object, &mut source, &options, &options) + .await + .expect("materialized self-copy"); + } else { + store + .put_object( + &bucket, + object, + &mut PutObjReader::from_vec(expected.clone()), + &ObjectOptions { + user_defined: new_metadata, + ..options.clone() + }, + ) + .await + .expect("overwrite transitioned null version"); + } + + let set = store.pools[0].get_disks_by_key(object); + let versions = set + .load_file_info_versions_exact(&bucket, object) + .await + .expect("read committed disk metadata") + .expect("replacement metadata exists"); + let free: Vec<_> = versions + .versions + .iter() + .chain(versions.free_versions.iter()) + .filter(|fi| fi.tier_free_version()) + .collect(); + assert_eq!(free.len(), 1, "{state:?}, suspended={suspended}, copy={self_copy}"); + assert_eq!(free[0].transitioned_objname, remote); + assert_eq!(free[0].transition_version_state, state); + assert!(backend.contains(&remote).await, "commit must not delete remote bytes before cleanup"); + let removed_before = backend.remove_count().await; + + // Restart before queue delivery. The new runtime must + // reconstruct ownership solely from the committed xl.meta. + let tier_config = ctx + .tier_config_mgr() + .read() + .await + .tiers + .get(tier) + .expect("tier configuration survives restart") + .clone_with_credentials(); + drop(set); + shutdown.cancel(); + drop(store); + drop(ctx); + (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite-restart", &[4])) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await; + { + let manager = ctx.tier_config_mgr(); + let mut manager = manager.write().await; + manager.tiers.insert(tier.to_string(), tier_config); + manager + .install_test_driver(tier, Box::new(backend.clone())) + .expect("rebind the same remote destination after restart"); + } + let set = store.pools[0].get_disks_by_key(object); + ExpiryState::resize_workers(1, Arc::clone(&store)).await; + let recovered = recover_tier_free_versions(Arc::clone(&store), 100, None, None) + .await + .expect("recover persisted cleanup owner"); + assert!(recovered.enqueued >= 1); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let versions = set + .load_file_info_versions_exact(&bucket, object) + .await + .expect("read cleanup progress") + .expect("new object must survive cleanup"); + if versions + .versions + .iter() + .chain(versions.free_versions.iter()) + .all(|fi| !fi.tier_free_version()) + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("cleanup must converge"); + assert!(!backend.contains(&remote).await); + assert_eq!(backend.remove_count().await, removed_before + 1, "one remote DELETE per owner"); + assert_eq!(backend.remove_versions().await.last(), Some(&(remote.clone(), version.to_string()))); + let mut reader = store + .get_object_reader(&bucket, object, None, HeaderMap::new(), &options) + .await + .expect("replacement remains readable"); + let mut actual = Vec::new(); + reader.stream.read_to_end(&mut actual).await.expect("read replacement bytes"); + assert_eq!(actual, expected); + let current = store + .get_object_info(&bucket, object, &options) + .await + .expect("replacement metadata"); + assert_eq!(current.user_defined.get("content-type").map(String::as_str), Some("text/plain")); + assert_eq!(current.user_defined.get("x-amz-meta-replacement").map(String::as_str), Some("kept")); + assert!(current.transitioned_object.status.is_empty()); + } + } + } + shutdown.cancel(); + } + #[cfg(feature = "test-util")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial_test::serial(storage_class_env)] diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index 9549afcd5..114c19634 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -480,7 +480,110 @@ impl FileMeta { Ok(()) } - pub fn add_version(&mut self, mut fi: FileInfo) -> Result<()> { + pub fn add_version(&mut self, fi: FileInfo) -> Result<()> { + if let Some(free_version) = self.overwritten_tier_free_version(&fi)? { + // The replacement and its cleanup owner must share one xl.meta + // commit. Keep the original intact if either insertion fails. + let mut next = self.clone(); + next.add_version_inner(fi)?; + next.add_version_filemata(free_version)?; + *self = next; + return Ok(()); + } + self.add_version_inner(fi) + } + + fn overwritten_tier_free_version(&self, fi: &FileInfo) -> Result> { + use rustfs_utils::http::{ + SUFFIX_TIER_FV_ID, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID, + SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, + get_consistent_bytes, get_consistent_str, has_internal_suffix, strip_internal_prefix_preserving_case, + }; + + if fi.version_id.is_some_and(|id| !id.is_nil()) || !contains_key_str(&fi.metadata, SUFFIX_TIER_FV_ID) { + return Ok(None); + } + let Some(existing) = self + .versions + .iter() + .find(|v| v.header.version_id.is_none_or(|id| id.is_nil())) + else { + return Ok(None); + }; + let old = existing.parse_version_meta()?; + let Some(mut object) = old.object else { + return Ok(None); + }; + let status = get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_STATUS); + if status.is_none() + && object + .meta_sys + .keys() + .any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_STATUS)) + { + // Empty status is a valid local object. The reader distinguishes + // it from conflicting aliases before the ordinary overwrite. + object.into_fileinfo(&fi.volume, &fi.name, false)?; + return Ok(None); + } + if status != Some(TRANSITION_COMPLETE.as_bytes()) { + return Ok(None); + } + // Reuse the reader's alias/state validation. A legacy empty remote + // version is valid and must not be mistaken for conflicting aliases. + object.into_fileinfo(&fi.volume, &fi.name, false)?; + if object + .meta_sys + .keys() + .any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_TIER_DESTINATION_ID)) + && get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID).is_none() + { + return Err(Error::FileCorrupt); + } + let transition_suffixes = [ + SUFFIX_TRANSITION_STATUS, + SUFFIX_TRANSITION_TIER, + SUFFIX_TRANSITION_TIER_DESTINATION_ID, + SUFFIX_TRANSITIONED_OBJECTNAME, + SUFFIX_TRANSITIONED_VERSION_ID, + SUFFIX_TRANSITIONED_VERSION_STATE, + ]; + let replacement = MetaObject::from(fi.clone()); + if transition_suffixes + .iter() + .all(|suffix| get_consistent_bytes(&object.meta_sys, suffix) == get_consistent_bytes(&replacement.meta_sys, suffix)) + { + return Ok(None); + } + let id = get_consistent_str(&fi.metadata, SUFFIX_TIER_FV_ID).ok_or(Error::FileCorrupt)?; + let id = Uuid::parse_str(id)?; + if id.is_nil() || self.versions.iter().any(|version| version.header.version_id == Some(id)) { + return Err(Error::FileCorrupt); + } + // The reader also accepts legacy key casing. Canonicalize only this + // cleanup source so init_free_version preserves every accepted field, + // including an explicitly empty unversioned remote version. + for suffix in transition_suffixes { + let value = object + .meta_sys + .iter() + .find(|(key, _)| { + strip_internal_prefix_preserving_case(key).is_some_and(|found| found.eq_ignore_ascii_case(suffix)) + }) + .map(|(_, value)| value.clone()); + if let Some(value) = value { + rustfs_utils::http::insert_bytes(&mut object.meta_sys, suffix, value); + } + } + let (free_version, created) = object.init_free_version(fi)?; + if !created { + return Err(Error::FileCorrupt); + } + Ok(Some(free_version)) + } + + fn add_version_inner(&mut self, mut fi: FileInfo) -> Result<()> { + rustfs_utils::http::remove_str(&mut fi.metadata, rustfs_utils::http::SUFFIX_TIER_FV_ID); // empty version_id means "null" (versioning disabled/suspended) if fi.version_id.is_none() { fi.version_id = Some(Uuid::nil()); @@ -1463,6 +1566,235 @@ mod test { }); } + fn tier_overwrite_fixture(state: crate::TransitionVersionState) -> (FileMeta, FileInfo) { + let mut source = FileInfo::new("object", 2, 2); + source.erasure.index = 1; + source.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixture timestamp")); + source.data_dir = Some(Uuid::new_v4()); + source.transition_status = TRANSITION_COMPLETE.to_string(); + source.transition_tier = "WARM".to_string(); + source.transitioned_objname = "remote/old-object".to_string(); + source.transition_version_state = state; + source.transition_version = match state { + crate::TransitionVersionState::Exact => Some("opaque-provider-version".to_string()), + crate::TransitionVersionState::SuspendedNull => Some("null".to_string()), + _ => None, + }; + rustfs_utils::http::insert_str( + &mut source.metadata, + rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + "ab".repeat(32), + ); + if state == crate::TransitionVersionState::KnownDisabled { + rustfs_utils::http::insert_str( + &mut source.metadata, + rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, + String::new(), + ); + } + let mut meta = FileMeta::new(); + meta.add_version(source.clone()).expect("seed transitioned null version"); + (meta, source) + } + + #[test] + fn tier_overwrite_preserves_exact_cleanup_owner_across_reload() { + use crate::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull, Unknown}; + use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX, SUFFIX_TIER_FV_ID}; + + for state in [Exact, KnownDisabled, SuspendedNull, Unknown] { + for inline in [false, true] { + let (mut meta, _) = tier_overwrite_fixture(state); + let old = meta.versions[0].parse_version_meta().expect("old metadata"); + let old = old.object.expect("old object"); + let id = Uuid::new_v4(); + let mut replacement = FileInfo::new("object", 2, 2); + replacement.version_id = inline.then_some(Uuid::nil()); + replacement.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_001).expect("fixture timestamp")); + replacement.data_dir = Some(Uuid::new_v4()); + replacement.size = 3; + if inline { + replacement.data = Some(Bytes::from_static(b"new")); + } + replacement.set_tier_free_version_id(&id.to_string()); + + meta.add_version(replacement.clone()) + .expect("replace transitioned null version"); + let bytes = meta.marshal_msg().expect("persist replacement and cleanup owner"); + let mut reopened = FileMeta::load(&bytes).expect("reopen committed metadata"); + assert_eq!(reopened.versions.len(), 2); + let (_, free) = reopened.find_version(Some(id)).expect("durable cleanup owner"); + assert!(free.free_version()); + let marker = free.delete_marker.expect("cleanup marker"); + for suffix in [ + rustfs_utils::http::SUFFIX_TRANSITION_TIER, + rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME, + rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, + rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE, + ] { + for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] { + let key = format!("{prefix}{suffix}"); + assert_eq!(marker.meta_sys.get(&key), old.meta_sys.get(&key), "{state:?}: {key}"); + } + } + let (_, current) = reopened.find_version(None).expect("replacement survives restart"); + let current = current.object.expect("replacement object"); + assert_eq!(current.size, 3); + assert!(!rustfs_utils::http::contains_key_bytes(¤t.meta_sys, SUFFIX_TIER_FV_ID)); + reopened + .add_version(replacement) + .expect("replaying replacement is idempotent"); + assert_eq!(reopened.versions.len(), 2); + } + } + } + + #[test] + fn tier_overwrite_rejects_cleanup_failure_without_mutating_source() { + for id in ["not-a-uuid".to_string(), Uuid::nil().to_string()] { + let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact); + let before = meta.clone(); + let mut replacement = FileInfo::new("object", 2, 2); + replacement.mod_time = Some(OffsetDateTime::now_utc()); + replacement.set_tier_free_version_id(&id); + assert!(meta.add_version(replacement).is_err()); + assert_eq!(meta, before, "invalid cleanup identity must preserve source"); + } + with_object_max_versions_for_test(1, || { + let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::KnownDisabled); + let before = meta.clone(); + let mut replacement = FileInfo::new("object", 2, 2); + replacement.mod_time = Some(OffsetDateTime::now_utc()); + replacement.data = Some(Bytes::from_static(b"new")); + replacement.set_tier_free_version_id(&Uuid::new_v4().to_string()); + assert_eq!( + meta.add_version(replacement) + .expect_err("cleanup owner exceeds version limit"), + Error::MaxVersionsExceeded + ); + assert_eq!(meta, before, "failed cleanup insertion must also preserve inline bytes"); + }); + } + + #[test] + fn tier_overwrite_allows_empty_transition_status_on_local_source() { + let mut source = FileInfo::new("object", 2, 2); + source.mod_time = Some(OffsetDateTime::now_utc()); + source.data = Some(Bytes::from_static(b"old")); + rustfs_utils::http::insert_str(&mut source.metadata, rustfs_utils::http::SUFFIX_TRANSITION_STATUS, String::new()); + let mut meta = FileMeta::new(); + meta.add_version(source) + .expect("seed readable local metadata with empty status"); + let mut replacement = FileInfo::new("object", 2, 2); + replacement.mod_time = Some(OffsetDateTime::now_utc()); + replacement.size = 3; + replacement.data = Some(Bytes::from_static(b"new")); + replacement.set_tier_free_version_id(&Uuid::new_v4().to_string()); + meta.add_version(replacement) + .expect("an ordinary overwrite must still succeed"); + assert_eq!(meta.versions.len(), 1); + let (_, current) = meta.find_version(None).expect("replacement remains visible"); + assert_eq!(current.object.expect("ordinary object").size, 3); + assert!(!meta.versions[0].header.free_version()); + } + + #[test] + fn tier_overwrite_preserves_legacy_metadata_casing() { + use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX}; + + for state in [ + crate::TransitionVersionState::Exact, + crate::TransitionVersionState::KnownDisabled, + ] { + let (mut meta, _) = tier_overwrite_fixture(state); + let mut source = meta.versions[0].parse_version_meta().expect("seeded source"); + let object = source.object.as_mut().expect("transitioned source"); + let expected = object.meta_sys.clone(); + object.meta_sys = object + .meta_sys + .drain() + .map(|(key, value)| (key.to_ascii_uppercase(), value)) + .collect(); + meta.versions[0] = FileMetaShallowVersion::try_from(source).expect("legacy key casing"); + let id = Uuid::new_v4(); + let mut replacement = FileInfo::new("object", 2, 2); + replacement.mod_time = Some(OffsetDateTime::now_utc()); + replacement.set_tier_free_version_id(&id.to_string()); + meta.add_version(replacement).expect("overwrite readable legacy source"); + let reopened = FileMeta::load(&meta.marshal_msg().expect("persist overwrite")).expect("reopen overwrite"); + let (_, owner) = reopened + .find_version(Some(id)) + .expect("legacy source must retain cleanup ownership"); + let marker = owner.delete_marker.expect("cleanup marker"); + for suffix in [ + rustfs_utils::http::SUFFIX_TRANSITION_TIER, + rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME, + rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, + rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE, + ] { + for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] { + let key = format!("{prefix}{suffix}"); + assert_eq!(marker.meta_sys.get(&key), expected.get(&key), "legacy {state:?}: {key}"); + } + } + } + } + + #[test] + fn tier_overwrite_rejects_conflicting_remote_metadata_aliases() { + use rustfs_utils::http::{ + MINIO_INTERNAL_PREFIX, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID, + SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, + }; + for suffix in [ + SUFFIX_TRANSITION_STATUS, + SUFFIX_TRANSITION_TIER, + SUFFIX_TRANSITION_TIER_DESTINATION_ID, + SUFFIX_TRANSITIONED_OBJECTNAME, + SUFFIX_TRANSITIONED_VERSION_ID, + SUFFIX_TRANSITIONED_VERSION_STATE, + ] { + let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact); + let mut old = meta.versions[0].parse_version_meta().expect("seeded source metadata"); + old.object + .as_mut() + .expect("transitioned source") + .meta_sys + .insert(format!("{MINIO_INTERNAL_PREFIX}{suffix}"), b"conflicting-value".to_vec()); + meta.versions[0] = FileMetaShallowVersion::try_from(old).expect("encode conflicting aliases"); + let before = meta.clone(); + let mut replacement = FileInfo::new("object", 2, 2); + replacement.mod_time = Some(OffsetDateTime::now_utc()); + replacement.set_tier_free_version_id(&Uuid::new_v4().to_string()); + assert_eq!( + meta.add_version(replacement) + .expect_err("ambiguous ownership must fail closed"), + Error::FileCorrupt + ); + assert_eq!(meta, before, "conflicting {suffix} must not erase the old remote tuple"); + } + } + + #[test] + fn tier_overwrite_keeps_retained_remote_and_versioned_copy_ownership() { + let (mut meta, mut source) = tier_overwrite_fixture(crate::TransitionVersionState::Exact); + source.set_tier_free_version_id(&Uuid::new_v4().to_string()); + meta.add_version(source.clone()) + .expect("restore retains the same remote owner"); + assert_eq!(meta.versions.len(), 1); + source.version_id = Some(Uuid::new_v4()); + source.transition_status.clear(); + source.transition_tier.clear(); + source.transitioned_objname.clear(); + source.transition_version = None; + source.transition_version_state = crate::TransitionVersionState::Unknown; + meta.add_version(source).expect("versioned write retains historical source"); + assert_eq!(meta.versions.len(), 2); + assert!(meta.versions.iter().all(|version| !version.header.free_version())); + } + #[test] fn add_version_filemata_uses_canonical_equal_time_order() { let mod_time = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid test timestamp"); diff --git a/docs/architecture/ilm-tiering-persistence-contracts.md b/docs/architecture/ilm-tiering-persistence-contracts.md index 74d34275d..4212c398e 100644 --- a/docs/architecture/ilm-tiering-persistence-contracts.md +++ b/docs/architecture/ilm-tiering-persistence-contracts.md @@ -30,11 +30,23 @@ These are approved-target invariants. A protocol's explicitly labeled current ex | Remote PUT is in flight or its response is unknown | Transition transaction | Only cleanup of its own canonical candidate, subject to the transaction recovery predicate | Durable transaction identity plus a known remote-version state; the approved target also requires expiry and durable takeover of the creator fence | | Local transition commit is complete | Exact transitioned version in `xl.meta` | No | Recovery finds the transaction's logical bucket/object/version and requires the complete recorded source identity (version ID, data directory, modification time, size, and ETag), `TRANSITION_COMPLETE`, and the same remote object, tier, and remote version before removing only the transaction record | | An ordinary delete removes that transitioned version | Hidden `xl.meta` free-version | Yes | Metadata quorum atomically removes the visible version and preserves its exact tier tuple in the free-version | +| PUT or materialized self-copy replaces a transitioned null version | Hidden `xl.meta` free-version | Yes, after replacement commit and complete physical reference checks | The coordinator supplies one cleanup UUID to every disk; replacement metadata and the old remote tuple are written in the same `xl.meta` commit | | A recursive prefix/delete-all operation cannot preserve per-object markers | v6 journal bound to an immutable single dispatch manifest or a chunk-parent-bound child manifest | Yes, but only after child/manifest completion and all-pool absence proof | `DispatchAuthorized`, exact local destructive mutation, every journal `Committed`, then child/manifest `Completed`; a chunk parent advances only after that child completion | | Tier configuration mutation, manual job, or decommission receipt | Intent/admission/copy proof only | No | These records gate configuration, scheduling, or migration; they never become remote-object cleanup owners | An old journal and a free-version can coexist during compatibility recovery. That coexistence is evidence of multiple possible owners, not permission to choose one: the journal path must retain its record until the version-specific recovery rule proves which owner is authoritative. +Null-version replacement uses the existing free-version format and recovery +worker. Failed metadata preparation preserves both the old version and its +inline bytes; rename rollback restores the complete previous metadata. Recovery +must retain cleanup while any physical replica still references the remote tuple, +including a minority version omitted by quorum merging, or any disk cannot be +checked. A successful replacement needs no in-memory queue receipt to survive +restart: the normal free-version sweep discovers its committed owner. Restores +that retain the same remote tuple and ordinary versioned writes retain their +existing ownership. Older binaries can read this format, but all writers and +cleanup workers need the overwrite fix before these guarantees cover the fleet. + ## Persisted record inventory All keys below are objects in the internal metadata bucket. The table gives the canonical target form. Transition-transaction runtime recovery extracts the final 32 hexadecimal characters and UUID while ignoring shard directories and accepting uppercase hex. Manual-job runtime recovery requires exactly two shards matching the filename prefix, but accepts an uppercase UUID when the shards use the same uppercase text; it then loads the lowercase canonical job by UUID. The decommission validator recomputes and rejects a noncanonical manual-job path, but currently inherits the weaker transition parser. Exact runtime canonical-path validation for both protocols is an approved target.