diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index c07a013de..6c653daa8 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -1492,6 +1492,7 @@ pub(in crate::set_disk) fn record_read_repair_dedup(reason: &'static str) { counter!("rustfs_heal_read_repair_dedup_total", "reason" => reason).increment(1); } +#[derive(Debug)] pub(in crate::set_disk) enum ReadRepairAdmissionOutcome { Response(HealAdmissionResult), Failed(String), @@ -3838,6 +3839,211 @@ pub(in crate::set_disk) struct RenameTailOutcome { pub(in crate::set_disk) cleanup: Vec, } +const EVENT_SET_DISK_RENAME_ROLLBACK: &str = "set_disk_rename_rollback"; + +#[derive(Debug, Clone, PartialEq, Eq)] +enum RenameRollbackOutcome { + NotAttempted(DiskError), + Succeeded, + Failed(DiskError), + Panicked, + Cancelled, +} + +impl RenameRollbackOutcome { + fn stage(&self) -> &'static str { + match self { + Self::NotAttempted(_) => "rename_failed", + Self::Succeeded => "undo_succeeded", + Self::Failed(_) => "undo_failed", + Self::Panicked => "undo_panicked", + Self::Cancelled => "undo_cancelled", + } + } + + fn failed(&self) -> bool { + matches!(self, Self::Failed(_) | Self::Panicked | Self::Cancelled) + } +} + +#[derive(Debug, Clone)] +struct RenameRollbackDiskOutcome { + disk_index: usize, + rollback_dir: Option, + outcome: RenameRollbackOutcome, +} + +#[derive(Debug)] +struct RenameRollbackReport { + disks: Vec, +} + +/// Shares rollback completion with the staging owner without replacing the +/// original disk/quorum error returned by the rename operation. +#[derive(Clone, Default)] +pub(in crate::set_disk) struct RenameRollbackReceipt(Arc>); + +impl RenameRollbackReceipt { + pub(in crate::set_disk) fn is_incomplete(&self) -> bool { + self.0 + .get() + .is_some_and(|report| report.disks.iter().any(|disk| disk.outcome.failed())) + } +} + +async fn inspect_incomplete_rename_rollback( + disks: &[Option], + bucket: &str, + object: &str, + submitter: ReadRepairAdmissionSubmitter, +) -> ReadRepairAdmissionOutcome { + let location = disks.iter().flatten().next().map(|disk| disk.get_disk_location()); + let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options( + bucket.to_string(), + Some(object.to_string()), + false, + Some(HealChannelPriority::High), + location.as_ref().and_then(|location| location.pool_idx), + location.as_ref().and_then(|location| location.set_idx), + ); + // A failed write's surviving minority is not an authoritative heal source. + // Request inspection only: MRF PartialWrite would schedule mutating repair. + request.dry_run = Some(true); + request.remove_corrupted = Some(false); + request.recreate_missing = Some(false); + request.update_parity = Some(false); + request.recursive = Some(false); + match tokio::time::timeout(Duration::from_secs(1), submitter(request)).await { + Ok(result) => result, + Err(_) => ReadRepairAdmissionOutcome::Failed("rollback inspection admission timed out".to_string()), + } +} + +fn rename_rollback_task_outcome( + result: std::result::Result, tokio::task::JoinError>, +) -> RenameRollbackOutcome { + match result { + Ok(Ok(())) => RenameRollbackOutcome::Succeeded, + Ok(Err(err)) => RenameRollbackOutcome::Failed(err), + Err(err) if err.is_panic() => RenameRollbackOutcome::Panicked, + Err(_) => RenameRollbackOutcome::Cancelled, + } +} + +async fn rollback_failed_rename( + disks: &[Option], + mut file_infos: Vec, + errs: &[Option], + rollback_dirs: &[Option], + dst: (&str, &str), + receipt: Option, +) { + let (bucket, object) = dst; + let mut outcomes = Vec::with_capacity(disks.len()); + let mut tasks = Vec::with_capacity(disks.len()); + for (disk_index, disk) in disks.iter().enumerate() { + let rollback_dir = rollback_dirs[disk_index]; + let outcome = match &errs[disk_index] { + Some(err) => RenameRollbackOutcome::NotAttempted(err.clone()), + None => RenameRollbackOutcome::Failed(DiskError::DiskNotFound), + }; + outcomes.push(RenameRollbackDiskOutcome { + disk_index, + rollback_dir, + outcome, + }); + if errs[disk_index].is_some() { + continue; + } + let Some(disk) = disk.clone() else { + continue; + }; + let fi = std::mem::take(&mut file_infos[disk_index]); + let bucket = bucket.to_string(); + let object = object.to_string(); + let task = tokio::spawn(async move { + #[allow(clippy::let_unit_value)] + let _task_guard = SetDisks::rename_fanout_task_guard(&object); + SetDisks::rename_fanout_barrier(&object, disk_index, rename_fanout_barrier_phase::ROLLBACK).await; + #[cfg(test)] + rollback_fault_injection::before_undo(&object, disk_index)?; + disk.delete_version( + &bucket, + &object, + fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir: rollback_dir, + ..Default::default() + }, + ) + .await + }); + tasks.push(async move { (disk_index, task.await) }); + } + for (disk_index, result) in join_all(tasks).await { + outcomes[disk_index].outcome = rename_rollback_task_outcome(result); + } + + let attempted = outcomes + .iter() + .filter(|disk| !matches!(disk.outcome, RenameRollbackOutcome::NotAttempted(_))) + .count(); + let failed = outcomes.iter().filter(|disk| disk.outcome.failed()).count(); + let succeeded = attempted - failed; + for disk in &outcomes { + counter!("rustfs_rename_rollback_disks_total", "stage" => disk.outcome.stage()).increment(1); + if disk.outcome.failed() { + let location = disks[disk.disk_index].as_ref().map(|disk| disk.get_disk_location()); + warn!( + event = EVENT_SET_DISK_RENAME_ROLLBACK, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + state = "recovery_required", + stage = disk.outcome.stage(), + pool_index = ?location.as_ref().and_then(|location| location.pool_idx), + set_index = ?location.as_ref().and_then(|location| location.set_idx), + disk_index = disk.disk_index, + bucket, + object, + rollback_dir = ?disk.rollback_dir, + outcome = ?disk.outcome, + attempted, + succeeded, + failed, + "rename rollback incomplete; preserve recovery material" + ); + } + } + if let Some(receipt) = receipt { + let _ = receipt.0.set(RenameRollbackReport { disks: outcomes }); + } + if failed > 0 { + let result = inspect_incomplete_rename_rollback(disks, bucket, object, send_read_repair_heal_request).await; + let admission = match &result { + ReadRepairAdmissionOutcome::Response(response) => response.result_label(), + ReadRepairAdmissionOutcome::Failed(_) => "failed", + }; + counter!("rustfs_rename_rollback_inspection_total", "admission" => admission).increment(1); + warn!( + event = EVENT_SET_DISK_RENAME_ROLLBACK, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + state = "recovery_required", + stage = "inspection_admission", + bucket, + object, + attempted, + succeeded, + failed, + admission, + outcome = ?result, + "rename rollback inspection requested; recovery remains incomplete" + ); + } +} + /// Options shared by the normal and early-ack rename fanouts. Keeping the /// quorum and optional scanner lease map together avoids widening either /// fanout helper's argument list while preserving the fence semantics. @@ -3845,6 +4051,7 @@ pub(in crate::set_disk) struct RenameDataFenceOptions<'a> { write_quorum: usize, scanner_publication_lease_tokens: Option<&'a HashMap>, scanner_publication_commit_scope: Option, + rollback_receipt: Option, } impl<'a> RenameDataFenceOptions<'a> { @@ -3856,9 +4063,15 @@ impl<'a> RenameDataFenceOptions<'a> { write_quorum, scanner_publication_lease_tokens, scanner_publication_commit_scope: None, + rollback_receipt: None, } } + pub(in crate::set_disk) fn with_rollback_receipt(mut self, receipt: RenameRollbackReceipt) -> Self { + self.rollback_receipt = Some(receipt); + self + } + pub(in crate::set_disk) fn with_publication_scope( mut self, scanner_publication_commit_scope: Option, @@ -4224,6 +4437,7 @@ impl SetDisks { write_quorum, scanner_publication_lease_tokens, scanner_publication_commit_scope: _scanner_publication_commit_scope, + rollback_receipt, } = fence_options; if let Some(file_info) = disks .iter() @@ -4405,36 +4619,15 @@ impl SetDisks { if !sent_commit { let ret_err = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum).unwrap_or(DiskError::Unexpected); - let mut rollbacks = Vec::new(); - let mut rollback_file_infos = file_infos; - for (i, err) in errs.iter().enumerate() { - if err.is_some() { - continue; - } - - if let Some(disk) = coordinator_disks[i].as_ref() { - let fi = std::mem::take(&mut rollback_file_infos[i]); - let old_data_dir = data_dirs[i]; - let disk = disk.clone(); - let dst_bucket = fanout_dst_bucket.clone(); - let dst_object = fanout_dst_object.clone(); - rollbacks.push(tokio::spawn(async move { - disk.delete_version( - &dst_bucket, - &dst_object, - fi, - false, - DeleteOptions { - undo_write: true, - old_data_dir, - ..Default::default() - }, - ) - .await - })); - } - } - let _ = join_all(rollbacks).await; + rollback_failed_rename( + &coordinator_disks, + file_infos, + &errs, + &data_dirs, + (&fanout_dst_bucket, &fanout_dst_object), + rollback_receipt, + ) + .await; if let Some(commit_tx) = commit_tx.take() { let _ = commit_tx.send(Err(ret_err)); } @@ -4582,6 +4775,7 @@ impl SetDisks { write_quorum, scanner_publication_lease_tokens, scanner_publication_commit_scope, + rollback_receipt, } = fence_options; if let Some(file_info) = disks .iter() @@ -4793,36 +4987,7 @@ impl SetDisks { ); } - let mut futures = Vec::with_capacity(disks.len()); if let Some(ret_err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - for (i, err) in errs.iter().enumerate() { - if err.is_some() { - continue; - } - - if let Some(disk) = disks[i].as_ref() { - let fi = std::mem::take(&mut file_infos[i]); - let old_data_dir = data_dirs[i]; - let disk = disk.clone(); - let dst_bucket = dst_bucket.clone(); - let dst_object = dst_object.clone(); - futures.push(tokio::spawn(async move { - disk.delete_version( - &dst_bucket, - &dst_object, - fi, - false, - DeleteOptions { - undo_write: true, - old_data_dir, - ..Default::default() - }, - ) - .await - })); - } - } - if issue3031_diag_enabled() { warn!( target: "rustfs_ecstore::set_disk", @@ -4838,23 +5003,7 @@ impl SetDisks { ); } - let undo_results = join_all(futures).await; - let undo_error_count = undo_results - .iter() - .filter(|result| match result { - Err(_) | Ok(Err(_)) => true, - Ok(Ok(_)) => false, - }) - .count(); - if undo_error_count > 0 { - warn!( - target: "rustfs_ecstore::set_disk", - dst_bucket = %dst_bucket, - dst_object = %dst_object, - undo_error_count, - "rename_data quorum rollback reported errors" - ); - } + rollback_failed_rename(disks, file_infos, &errs, &data_dirs, (&dst_bucket, &dst_object), rollback_receipt).await; return Err(ret_err); } @@ -6800,6 +6949,57 @@ pub(in crate::set_disk) mod rename_fault_injection { } } +#[cfg(test)] +pub(in crate::set_disk) mod rollback_fault_injection { + use super::DiskError; + use std::{ + collections::HashMap, + sync::{Mutex, OnceLock}, + }; + + #[derive(Clone, Copy, Debug)] + pub(in crate::set_disk) enum Fault { + Io, + Panic, + } + + fn registry() -> &'static Mutex> { + static REGISTRY: OnceLock>> = OnceLock::new(); + REGISTRY.get_or_init(Mutex::default) + } + + pub(in crate::set_disk) struct Guard(String); + + impl Drop for Guard { + fn drop(&mut self) { + if let Ok(mut registry) = registry().lock() { + registry.remove(&self.0); + } + } + } + + pub(in crate::set_disk) fn arm(object: &str, disk_index: usize, fault: Fault) -> Guard { + registry() + .lock() + .expect("rollback registry should not poison") + .insert(object.to_string(), (disk_index, fault)); + Guard(object.to_string()) + } + + pub(super) fn before_undo(object: &str, disk_index: usize) -> Result<(), DiskError> { + let fault = registry() + .lock() + .expect("rollback registry should not poison") + .get(object) + .copied(); + match fault { + Some((target, Fault::Io)) if target == disk_index => Err(DiskError::FaultyDisk), + Some((target, Fault::Panic)) if target == disk_index => panic!("injected rollback panic"), + _ => Ok(()), + } + } +} + /// Test-only per-disk call counters for the metadata fan-out (backlog#1325, /// serving the RPC-count assertions of #1309 / #1314 / #1315). /// @@ -6911,6 +7111,7 @@ pub(in crate::set_disk) mod rename_fanout_barrier_phase { pub const RENAME: &str = "rename"; /// The per-disk old-data-dir cleanup phase of the commit fan-out. pub const CLEANUP: &str = "cleanup"; + pub const ROLLBACK: &str = "rollback"; /// The per-disk `read_version` phase of metadata read fan-out. #[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")] pub const READ_VERSION: &str = "read_version"; @@ -10422,6 +10623,309 @@ mod tests { .await; } + #[tokio::test] + async fn rename_rollback_incomplete_inspection_rejection_is_not_recovery() { + fn reject_inspection(request: rustfs_heal_contracts::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture { + assert_eq!(request.bucket, "rollback-inspection"); + assert_eq!(request.object_prefix.as_deref(), Some("object")); + assert_eq!(request.dry_run, Some(true), "failed minority must never become a mutating heal source"); + assert_eq!(request.remove_corrupted, Some(false)); + assert_eq!(request.recreate_missing, Some(false)); + assert_eq!(request.recursive, Some(false)); + Box::pin(async { ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Full) }) + } + let result = inspect_incomplete_rename_rollback(&[], "rollback-inspection", "object", reject_inspection).await; + assert!(matches!(result, ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Full))); + } + + #[tokio::test] + async fn rename_rollback_incomplete_cancelled_task_is_not_success() { + let task = tokio::spawn(std::future::pending::>()); + task.abort(); + assert_eq!(rename_rollback_task_outcome(task.await), RenameRollbackOutcome::Cancelled); + } + + #[tokio::test] + #[serial_test::serial(capacity_dirty_scope)] + async fn rename_rollback_incomplete_matches_early_ack_and_full_wait_after_reopen() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + const DISKS: usize = 4; + const WRITE_QUORUM: usize = 3; + for overwrite in [false, true] { + for success_count in [0, WRITE_QUORUM - 1, WRITE_QUORUM] { + for fault in [ + None, + Some(rollback_fault_injection::Fault::Io), + Some(rollback_fault_injection::Fault::Panic), + ] { + let mut previous = None; + for early_ack in [false, true] { + let bucket = "rename-rollback-matrix"; + let object = format!("object-{overwrite}-{success_count}-{fault:?}-{early_ack}"); + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + if overwrite { + let mut old = metadata_test_fileinfo(&object); + old.mod_time = Some(OffsetDateTime::now_utc()); + old.data = Some(Bytes::from_static(b"old-inline-body")); + old.metadata.insert("etag".to_string(), "old-etag".to_string()); + for disk in disks.iter().flatten() { + disk.write_metadata(bucket, bucket, &object, old.clone()) + .await + .expect("old version must be staged"); + } + } + let _rename_fault = + rename_fault_injection::fail_rename_on(&object, &(success_count..DISKS).collect::>()); + let _undo_fault = fault.map(|fault| rollback_fault_injection::arm(&object, 0, fault)); + let receipt = RenameRollbackReceipt::default(); + let result = SetDisks::rename_data_owned_with_fence( + &disks, + (RUSTFS_META_TMP_BUCKET, "source"), + rename_commit_fileinfos(&object, DISKS, "new-etag"), + (bucket, &object), + early_ack, + RenameDataFenceOptions::new(WRITE_QUORUM, None).with_rollback_receipt(receipt.clone()), + ) + .await; + let actual_error = match result { + Ok(commit) => { + assert_eq!(success_count, WRITE_QUORUM); + if let Some(tail) = commit.tail_drain { + tail.await + .expect("committed tail should join") + .expect("committed tail should converge"); + } + assert!( + receipt.0.get().is_none(), + "a quorum commit must not enter rollback even when undo faults are armed" + ); + None + } + Err(err) => { + assert!(success_count < WRITE_QUORUM); + let report = receipt + .0 + .get() + .expect("both failure paths must publish per-disk rollback evidence"); + assert_eq!(report.disks.len(), DISKS); + for (idx, outcome) in report.disks.iter().enumerate() { + assert_eq!(outcome.disk_index, idx); + let expected = if idx >= success_count { + RenameRollbackOutcome::NotAttempted(DiskError::other( + "injected rename failure (test-only)", + )) + } else if idx == 0 { + match fault { + Some(rollback_fault_injection::Fault::Io) => { + RenameRollbackOutcome::Failed(DiskError::FaultyDisk) + } + Some(rollback_fault_injection::Fault::Panic) => RenameRollbackOutcome::Panicked, + None => RenameRollbackOutcome::Succeeded, + } + } else { + RenameRollbackOutcome::Succeeded + }; + assert_eq!(outcome.outcome, expected); + if overwrite && idx < success_count { + let backup = dirs[idx] + .path() + .join(bucket) + .join(&object) + .join(outcome.rollback_dir.expect("overwrite needs rollback dir").to_string()) + .join(STORAGE_FORMAT_FILE_BACKUP); + assert_eq!( + backup.exists(), + outcome.outcome.failed(), + "failed undo must retain its only old-version backup" + ); + } + } + assert_eq!(receipt.is_incomplete(), success_count > 0 && fault.is_some()); + Some(err) + } + }; + if let Some(expected_error) = previous.as_ref() { + assert_eq!( + &actual_error, expected_error, + "early ACK must preserve the original full-wait quorum error" + ); + } + previous = Some(actual_error); + for (idx, dir) in dirs.iter().enumerate() { + let reopened = reopen_local_disk(dir).await; + let read = reopened + .read_version( + "", + bucket, + &object, + "", + &ReadOptions { + read_data: true, + ..Default::default() + }, + ) + .await; + let keeps_new = + idx < success_count && (success_count == WRITE_QUORUM || (idx == 0 && fault.is_some())); + if keeps_new || overwrite { + let stored = read.expect("old or committed version must survive reopen"); + assert_eq!( + stored.metadata.get("etag").map(String::as_str), + Some(if keeps_new { "new-etag" } else { "old-etag" }) + ); + assert_eq!( + stored.data.as_deref(), + Some(if keeps_new { + b"inline-body".as_slice() + } else { + b"old-inline-body".as_slice() + }) + ); + } else { + assert!( + matches!(read, Err(DiskError::FileNotFound | DiskError::FileVersionNotFound)), + "fresh rollback must not expose data: {read:?}" + ); + } + } + } + } + } + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial(capacity_dirty_scope)] + async fn rename_rollback_incomplete_preserves_overwrite_data_dirs_and_staging() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + for early_ack in [false, true] { + let bucket = "rollback-data-dirs"; + let object = format!("object-{early_ack}"); + let (dirs, disks) = call_counter_local_disks(bucket, 4).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + let old_data_dir = Uuid::new_v4(); + let new_data_dir = Uuid::new_v4(); + let mut old = metadata_test_fileinfo(&object); + old.data_dir = Some(old_data_dir); + old.mod_time = Some(OffsetDateTime::now_utc()); + let mut infos = Vec::new(); + for (idx, disk) in disks.iter().enumerate() { + let disk = disk.as_ref().expect("fixture disk should be present"); + disk.write_metadata(bucket, bucket, &object, old.clone()) + .await + .expect("old metadata should be staged"); + let old_dir = dirs[idx].path().join(bucket).join(&object).join(old_data_dir.to_string()); + tokio::fs::create_dir_all(&old_dir) + .await + .expect("old data directory should exist"); + tokio::fs::write(old_dir.join("part.1"), b"old-data") + .await + .expect("old shard should exist"); + let source = dirs[idx] + .path() + .join(RUSTFS_META_TMP_BUCKET) + .join("source") + .join(new_data_dir.to_string()); + tokio::fs::create_dir_all(&source) + .await + .expect("new data directory should be staged"); + tokio::fs::write(source.join("part.1"), b"new-data") + .await + .expect("new shard should be staged"); + let mut fi = metadata_test_fileinfo(&object); + fi.data_dir = Some(new_data_dir); + fi.erasure.index = idx + 1; + fi.mod_time = Some(OffsetDateTime::now_utc()); + infos.push(fi); + } + let _rename_fault = rename_fault_injection::fail_rename_on(&object, &[2, 3]); + let _undo_fault = rollback_fault_injection::arm(&object, 0, rollback_fault_injection::Fault::Io); + let receipt = RenameRollbackReceipt::default(); + assert!( + SetDisks::rename_data_owned_with_fence( + &disks, + (RUSTFS_META_TMP_BUCKET, "source"), + infos, + (bucket, &object), + early_ack, + RenameDataFenceOptions::new(3, None).with_rollback_receipt(receipt.clone()), + ) + .await + .is_err() + ); + assert!(receipt.is_incomplete()); + for (idx, dir) in dirs.iter().enumerate() { + let root = dir.path().join(bucket).join(&object); + assert_eq!( + tokio::fs::read(root.join(old_data_dir.to_string()).join("part.1")) + .await + .expect("old data must survive failed overwrite"), + b"old-data" + ); + let backup = root.join(old_data_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP); + assert_eq!(backup.exists(), idx == 0, "only the failed undo retains its backup"); + if idx >= 2 { + let staged = dir + .path() + .join(RUSTFS_META_TMP_BUCKET) + .join("source") + .join(new_data_dir.to_string()) + .join("part.1"); + assert_eq!( + tokio::fs::read(staged) + .await + .expect("failed-write staging must remain available"), + b"new-data" + ); + } + let reopened = reopen_local_disk(dir).await; + let stored = reopened + .read_version("", bucket, &object, "", &ReadOptions::default()) + .await + .expect("metadata should survive reopen"); + assert_eq!(stored.data_dir, Some(if idx == 0 { new_data_dir } else { old_data_dir })); + } + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial(capacity_dirty_scope)] + async fn rename_rollback_incomplete_receipt_waits_for_undo_barrier() { + let bucket = "rename-rollback-barrier"; + let object = "rollback-barrier-object"; + let (dirs, disks) = call_counter_local_disks(bucket, 4).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + let _rename_fault = rename_fault_injection::fail_rename_on(object, &[2, 3]); + let _undo_fault = rollback_fault_injection::arm(object, 0, rollback_fault_injection::Fault::Io); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier_phase::ROLLBACK); + let receipt = RenameRollbackReceipt::default(); + let mut rename = Box::pin(SetDisks::rename_data_owned_with_fence( + &disks, + (RUSTFS_META_TMP_BUCKET, "source"), + rename_commit_fileinfos(object, 4, "new-etag"), + (bucket, object), + false, + RenameDataFenceOptions::new(3, None).with_rollback_receipt(receipt.clone()), + )); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + tokio::select! { + () = barrier.wait_until_paused() => {} + _ = rename.as_mut() => panic!("rename returned before the armed rollback barrier"), + } + }) + .await + .expect("undo must reach its disk barrier"); + assert!(receipt.0.get().is_none(), "pending undo must not be recorded as success"); + barrier.release(); + assert!(rename.await.is_err()); + assert!(receipt.is_incomplete(), "drained undo failure must survive in the receipt"); + } + #[tokio::test] #[serial_test::serial(capacity_dirty_scope)] async fn rename_data_early_ack_strict_quorum_failure_rolls_back_fresh_after_reopen() {