diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 1840058ae..3c0d2ff45 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -3732,6 +3732,103 @@ mod transition_commit_failure_tests { assert_eq!(restored, payload); } + #[tokio::test] + #[serial_test::serial] + async fn partial_local_commit_failure_rolls_back_applied_disks_and_preserves_remote_candidate() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "transition-partial-rollback-bucket"; + let object = "object.bin"; + let payload = b"partial local commit failure must roll back applied disks".repeat(1024); + for disk in &disk_stores { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let mut reader = PutObjReader::from_vec(payload.clone()); + let original = set_disks + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("source object should be written"); + + 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; + let opts = ObjectOptions { + no_lock: true, + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name, + etag: original.etag.clone().unwrap_or_default(), + ..Default::default() + }, + version_id: original.version_id.map(|version| version.to_string()), + mod_time: original.mod_time, + ..Default::default() + }; + + let barrier = TransitionCommitBarrier::install_after_lease_check(bucket, object); + let transition_set = Arc::clone(&set_disks); + let transition = tokio::spawn(async move { transition_set.transition_object(bucket, object, &opts).await }); + barrier.wait_until_paused().await; + assert_eq!(backend.put_count().await, 1, "remote candidate must exist before local commit"); + + let saved_disks = { + let mut disks = set_disks.disks.write().await; + let saved = disks.clone(); + for disk in disks.iter_mut().skip(2) { + *disk = None; + } + saved + }; + barrier.release(); + let result = transition.await.expect("transition task should not panic"); + + result.expect_err("partial local commit must fail when only two of four disks are writable"); + let fi = set_disks + .get_object_fileinfo( + bucket, + object, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + true, + false, + ) + .await + .expect("rollback should keep the source metadata readable on applied disks") + .0; + assert_ne!( + fi.transition_status, TRANSITION_COMPLETE, + "rollback must not leave the applied disks marked as transitioned" + ); + let mut restored = Vec::new(); + set_disks + .get_object_reader( + bucket, + object, + None, + HeaderMap::new(), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("rollback should keep the source object readable on applied disks") + .stream + .read_to_end(&mut restored) + .await + .expect("the local source body should drain after rollback"); + assert_eq!(restored, payload); + + *set_disks.disks.write().await = saved_disks; + assert_eq!( + backend.remove_count().await, + 0, + "ambiguous local commit failure must retain the remote candidate for scanner reconciliation" + ); + assert_eq!(backend.object_count().await, 1); + } + #[tokio::test] #[serial_test::serial] async fn production_transition_replacement_revokes_commit_and_cleans_with_old_driver() {