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
This commit is contained in:
cxymds
2026-09-10 20:14:07 +08:00
committed by GitHub
parent ffb18979f8
commit 00aeb12914
9 changed files with 904 additions and 8 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5
sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a
+1 -1
View File
@@ -1 +1 @@
sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0
sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9
+77 -2
View File
@@ -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
@@ -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]
@@ -3076,6 +3076,23 @@ impl SetDisks {
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
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<Option<rustfs_filemeta::FileInfoVersions>> {
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<Option<rustfs_filemeta::FileInfoVersions>> {
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";
+103
View File
@@ -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;
+218
View File
@@ -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)]
+333 -1
View File
@@ -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<Option<FileMetaVersion>> {
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(&current.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");
@@ -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.