diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 1f9e7e86f..a1d88b1b9 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -4043,6 +4043,7 @@ mod transition_upload_integrity_tests { use crate::services::tier::test_util::register_mock_tier; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; use http::HeaderMap; + use std::path::PathBuf; async fn assert_local_source_intact(set_disks: &Arc, bucket: &str, object: &str, payload: &[u8]) { let mut restored = Vec::new(); @@ -4113,6 +4114,56 @@ mod transition_upload_integrity_tests { } } + async fn corrupt_shard_at(path: PathBuf, position: ShardCorruptionPosition) { + let mut bytes = tokio::fs::read(&path).await.expect("committed shard should be readable"); + assert!(!bytes.is_empty(), "committed shard should not be empty"); + let offset = match position { + ShardCorruptionPosition::First => 0, + ShardCorruptionPosition::Middle => bytes.len() / 2, + ShardCorruptionPosition::Last => bytes.len() - 1, + }; + bytes[offset] ^= 0xff; + tokio::fs::write(&path, bytes) + .await + .expect("committed shard corruption should be written"); + } + + #[derive(Clone, Copy, Debug)] + enum ShardCorruptionPosition { + First, + Middle, + Last, + } + + impl ShardCorruptionPosition { + fn label(self) -> &'static str { + match self { + ShardCorruptionPosition::First => "first", + ShardCorruptionPosition::Middle => "middle", + ShardCorruptionPosition::Last => "last", + } + } + } + + async fn corrupt_beyond_read_quorum( + temp_dirs: &[tempfile::TempDir], + bucket: &str, + object: &str, + data_dir: Uuid, + parity_blocks: usize, + position: ShardCorruptionPosition, + ) { + for temp_dir in temp_dirs.iter().take(parity_blocks + 1) { + let shard = temp_dir + .path() + .join(bucket) + .join(object) + .join(data_dir.to_string()) + .join("part.1"); + corrupt_shard_at(shard, position).await; + } + } + #[tokio::test] #[serial_test::serial] async fn partial_remote_acceptance_cleans_exact_candidate_and_preserves_source() { @@ -4142,6 +4193,98 @@ mod transition_upload_integrity_tests { assert_local_source_intact(&set_disks, bucket, object, &payload).await; } + #[tokio::test] + #[serial_test::serial] + async fn real_bitrot_producer_failures_do_not_commit_transition() { + let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await; + + for position in [ + ShardCorruptionPosition::First, + ShardCorruptionPosition::Middle, + ShardCorruptionPosition::Last, + ] { + let bucket = format!("transition-real-bitrot-{}", position.label()); + let object = format!("{}-corrupt.bin", position.label()); + let payload = vec![0x41; 2 * 1024 * 1024]; + let original = write_source(&set_disks, &disk_stores, &bucket, &object, &payload).await; + let (source, _, _) = set_disks + .get_object_fileinfo( + &bucket, + &object, + &ObjectOptions { + no_lock: true, + metadata_cache_safe: false, + ..Default::default() + }, + true, + false, + ) + .await + .expect("source metadata should be available before shard corruption"); + let data_dir = source.data_dir.expect("source object should have a data directory"); + + corrupt_beyond_read_quorum(&temp_dirs, &bucket, &object, data_dir, source.erasure.parity_blocks, position).await; + + let error = set_disks + .transition_object(&bucket, &object, &transition_options(&original, tier_name.clone())) + .await + .expect_err("producer bitrot failure must not commit transition metadata"); + assert!( + matches!( + error, + StorageError::FileCorrupt + | StorageError::ErasureReadQuorum + | StorageError::InsufficientReadQuorum(_, _) + | StorageError::LessData + | StorageError::Io(_) + ), + "{position:?}: unexpected transition producer error: {error:?}" + ); + let (after, _, _) = set_disks + .get_object_fileinfo( + &bucket, + &object, + &ObjectOptions { + no_lock: true, + metadata_cache_safe: false, + ..Default::default() + }, + true, + false, + ) + .await + .expect("failed transition must leave metadata readable"); + assert_eq!( + after.data_dir, + Some(data_dir), + "{position:?}: transition must not release the source data dir" + ); + assert_ne!( + after.transition_status, TRANSITION_COMPLETE, + "{position:?}: transition status must remain incomplete after producer bitrot failure" + ); + } + + assert_eq!( + backend.put_count().await, + 3, + "each bitrot case should create one remote cleanup candidate" + ); + assert_eq!(backend.remove_count().await, 3, "each bitrot candidate must be removed before returning"); + assert_eq!( + backend.remove_versions().await, + backend.put_versions().await, + "bitrot cleanup must target the exact remote version returned by PUT" + ); + assert_eq!( + backend.object_count().await, + 0, + "bitrot producer failures must not leave remote candidates" + ); + } + #[tokio::test] #[serial_test::serial] async fn opaque_remote_version_is_cleaned_before_parse_failure() {