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