diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 6c653daa8..d07651ba8 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -3844,6 +3844,7 @@ const EVENT_SET_DISK_RENAME_ROLLBACK: &str = "set_disk_rename_rollback"; #[derive(Debug, Clone, PartialEq, Eq)] enum RenameRollbackOutcome { NotAttempted(DiskError), + Indeterminate(DiskError), Succeeded, Failed(DiskError), Panicked, @@ -3853,7 +3854,8 @@ enum RenameRollbackOutcome { impl RenameRollbackOutcome { fn stage(&self) -> &'static str { match self { - Self::NotAttempted(_) => "rename_failed", + Self::NotAttempted(_) => "rename_not_dispatched", + Self::Indeterminate(_) => "rename_indeterminate", Self::Succeeded => "undo_succeeded", Self::Failed(_) => "undo_failed", Self::Panicked => "undo_panicked", @@ -3861,8 +3863,12 @@ impl RenameRollbackOutcome { } } - fn failed(&self) -> bool { - matches!(self, Self::Failed(_) | Self::Panicked | Self::Cancelled) + fn undo_attempted(&self) -> bool { + matches!(self, Self::Succeeded | Self::Failed(_) | Self::Panicked | Self::Cancelled) + } + + fn needs_recovery(&self) -> bool { + matches!(self, Self::Indeterminate(_) | Self::Failed(_) | Self::Panicked | Self::Cancelled) } } @@ -3887,7 +3893,7 @@ 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())) + .is_some_and(|report| report.disks.iter().any(|disk| disk.outcome.needs_recovery())) } } @@ -3932,69 +3938,122 @@ fn rename_rollback_task_outcome( async fn rollback_failed_rename( disks: &[Option], - mut file_infos: Vec, + file_infos: Vec, errs: &[Option], + dispatched: &[bool], 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 owned_disks = disks.to_vec(); + let owned_errs = errs.to_vec(); + let owned_dispatched = dispatched.to_vec(); + let owned_dirs = rollback_dirs.to_vec(); + let owned_dst = (dst.0.to_string(), dst.1.to_string()); + let coordinator_failure_receipt = receipt.clone(); + // Own both undo mutations and their accounting: a cancelled requester must + // not leave detached disk tasks without the recovery evidence they produce. + let rollback = tokio::spawn(async move { + let disks = owned_disks.as_slice(); + let errs = owned_errs.as_slice(); + let dispatched = owned_dispatched.as_slice(); + let rollback_dirs = owned_dirs.as_slice(); + let dst = (owned_dst.0.as_str(), owned_dst.1.as_str()); + let mut file_infos = file_infos; - let attempted = outcomes + 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) if dispatched[disk_index] => RenameRollbackOutcome::Indeterminate(err.clone()), + 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); + } + + record_rename_rollback_outcomes(disks, outcomes, dst, receipt).await; + }); + if rollback.await.is_err() { + record_indeterminate_rename(disks, dst, coordinator_failure_receipt).await; + } +} + +async fn record_indeterminate_rename(disks: &[Option], dst: (&str, &str), receipt: Option) { + let outcomes = disks .iter() - .filter(|disk| !matches!(disk.outcome, RenameRollbackOutcome::NotAttempted(_))) + .enumerate() + .map(|(disk_index, disk)| RenameRollbackDiskOutcome { + disk_index, + rollback_dir: None, + outcome: if disk.is_some() { + RenameRollbackOutcome::Indeterminate(DiskError::Unexpected) + } else { + RenameRollbackOutcome::NotAttempted(DiskError::DiskNotFound) + }, + }) + .collect(); + record_rename_rollback_outcomes(disks, outcomes, dst, receipt).await; +} + +async fn record_rename_rollback_outcomes( + disks: &[Option], + outcomes: Vec, + dst: (&str, &str), + receipt: Option, +) { + let (bucket, object) = dst; + let attempted = outcomes.iter().filter(|disk| disk.outcome.undo_attempted()).count(); + let failed = outcomes + .iter() + .filter(|disk| disk.outcome.undo_attempted() && disk.outcome.needs_recovery()) + .count(); + let indeterminate = outcomes + .iter() + .filter(|disk| matches!(disk.outcome, RenameRollbackOutcome::Indeterminate(_))) .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() { + if disk.outcome.needs_recovery() { let location = disks[disk.disk_index].as_ref().map(|disk| disk.get_disk_location()); warn!( event = EVENT_SET_DISK_RENAME_ROLLBACK, @@ -4012,6 +4071,7 @@ async fn rollback_failed_rename( attempted, succeeded, failed, + indeterminate, "rename rollback incomplete; preserve recovery material" ); } @@ -4019,7 +4079,7 @@ async fn rollback_failed_rename( if let Some(receipt) = receipt { let _ = receipt.0.set(RenameRollbackReport { disks: outcomes }); } - if failed > 0 { + if failed > 0 || indeterminate > 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(), @@ -4037,6 +4097,7 @@ async fn rollback_failed_rename( attempted, succeeded, failed, + indeterminate, admission, outcome = ?result, "rename rollback inspection requested; recovery remains incomplete" @@ -4461,6 +4522,7 @@ impl SetDisks { let dst_object = Arc::new(dst_object.to_string()); let (commit_tx, commit_rx) = tokio::sync::oneshot::channel(); + let coordinator_failure_receipt = rollback_receipt.clone(); let tail_drain = tokio::spawn({ let fanout_src_bucket = src_bucket.clone(); let fanout_src_object = src_object.clone(); @@ -4483,7 +4545,8 @@ impl SetDisks { let file_info = file_info.clone(); let successful_rename_completion_rank = successful_rename_completion_rank.clone(); tasks.spawn(async move { - let result = std::panic::AssertUnwindSafe(async move { + let mut dispatched = false; + let result = std::panic::AssertUnwindSafe(async { #[allow(clippy::let_unit_value)] let _fanout_task_guard = Self::rename_fanout_task_guard(&dst_object); @@ -4508,6 +4571,7 @@ impl SetDisks { } let disk_wait_started = rustfs_io_metrics::put_stage_timer(); + dispatched = true; let result = disk .rename_data_borrowed_with_fence( &src_bucket, @@ -4518,6 +4582,10 @@ impl SetDisks { scanner_publication_lease_token, ) .await; + #[cfg(test)] + if result.is_ok() { + rollback_fault_injection::after_rename(&dst_object, i)?; + } if let Some(disk_wait_started) = disk_wait_started { let duration_ms = disk_wait_started.elapsed().as_secs_f64() * 1000.0; rustfs_io_metrics::record_put_object_stage_duration( @@ -4543,7 +4611,7 @@ impl SetDisks { }) .catch_unwind() .await; - (i, result) + (i, dispatched, result) }); } @@ -4553,6 +4621,8 @@ impl SetDisks { let mut fanout_panic = 0usize; let mut results_seen = 0usize; let mut errs = vec![Some(DiskError::DiskNotFound); disk_count]; + // Missing task results cannot prove that a disk mutation never ran. + let mut dispatched = vec![true; disk_count]; let mut disk_versions = vec![None; disk_count]; let mut data_dirs = vec![None; disk_count]; let mut cleanup_data_dirs = vec![None; disk_count]; @@ -4563,7 +4633,8 @@ impl SetDisks { while let Some(joined) = tasks.join_next().await { results_seen += 1; match joined { - Ok((idx, Ok(Ok(res)))) => { + Ok((idx, was_dispatched, Ok(Ok(res)))) => { + dispatched[idx] = was_dispatched; data_dirs[idx] = res.rollback_data_dir.or(res.old_data_dir); cleanup_data_dirs[idx] = res.cleanup_data_dir; disk_versions[idx] = res.sign; @@ -4571,10 +4642,12 @@ impl SetDisks { errs[idx] = None; success_count += 1; } - Ok((idx, Ok(Err(err)))) => { + Ok((idx, was_dispatched, Ok(Err(err)))) => { + dispatched[idx] = was_dispatched; errs[idx] = Some(err); } - Ok((idx, Err(_))) => { + Ok((idx, was_dispatched, Err(_))) => { + dispatched[idx] = was_dispatched; errs[idx] = Some(DiskError::Unexpected); fanout_panic += 1; } @@ -4604,6 +4677,8 @@ impl SetDisks { } } + #[cfg(test)] + rollback_fault_injection::after_fanout(&fanout_dst_object); if rustfs_io_metrics::put_stage_metrics_enabled() { let fanout_success = errs.iter().filter(|err| err.is_none()).count(); let fanout_error = errs.len().saturating_sub(fanout_success + fanout_panic); @@ -4623,6 +4698,7 @@ impl SetDisks { &coordinator_disks, file_infos, &errs, + &dispatched, &data_dirs, (&fanout_dst_bucket, &fanout_dst_object), rollback_receipt, @@ -4717,7 +4793,13 @@ impl SetDisks { }); let quorum_wait_started = rustfs_io_metrics::put_stage_timer(); - let commit = commit_rx.await.map_err(|_| DiskError::Unexpected)?; + let commit = match commit_rx.await { + Ok(commit) => commit, + Err(_) => { + record_indeterminate_rename(disks, (&dst_bucket, &dst_object), coordinator_failure_receipt).await; + return Err(DiskError::Unexpected); + } + }; rustfs_io_metrics::record_put_object_stage_duration_from( rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_QUORUM_WAIT, quorum_wait_started, @@ -4831,81 +4913,95 @@ impl SetDisks { let successful_rename_completion_rank = successful_rename_completion_rank.clone(); let publication_scope = scanner_publication_commit_scope.clone(); - std::panic::AssertUnwindSafe(async move { - // Test-only introspection guard: counts this operation as - // in-flight for the whole body. Compiles to `()` in production. - #[allow(clippy::let_unit_value)] - let _fanout_task_guard = Self::rename_fanout_task_guard(&dst_object); + async move { + let mut dispatched = false; + let result = std::panic::AssertUnwindSafe(async { + // Test-only introspection guard: counts this operation as + // in-flight for the whole body. Compiles to `()` in production. + #[allow(clippy::let_unit_value)] + let _fanout_task_guard = Self::rename_fanout_task_guard(&dst_object); - let Some(disk) = disk else { - return Err(DiskError::DiskNotFound); - }; - - let is_delete_marker = file_info.is_canonical_delete_marker(); - let mut local_file_info; - let file_info = if file_info.erasure.index == 0 { - local_file_info = file_info.clone(); - local_file_info.erasure.index = i + 1; - &local_file_info - } else { - file_info - }; - if file_info.erasure.index == 0 || (!is_delete_marker && !file_info.has_valid_erasure_geometry()) { - return Err(DiskError::FileCorrupt); - } - - // Test-only awaitable pause point right before the disk rename. - // A no-op immediately-ready future in production. - Self::rename_fanout_barrier(&dst_object, i, rename_fanout_barrier_phase::RENAME).await; - - if let Some(err) = Self::rename_injected_error(&dst_object, i) { - return Err(err); - } - - if let Some(scope) = publication_scope.as_ref() - && !scope.can_commit() - { - let _ = scope.mark_indeterminate(); - return Err(DiskError::other("scanner publication commit scope deadline or cancellation reached")); - } - - let disk_wait_started = rustfs_io_metrics::put_stage_timer(); - let result = disk - .rename_data_borrowed_with_fence( - &src_bucket, - &src_object, - file_info, - &dst_bucket, - &dst_object, - scanner_publication_lease_token, - ) - .await; - if let Some(disk_wait_started) = disk_wait_started { - let duration_ms = disk_wait_started.elapsed().as_secs_f64() * 1000.0; - rustfs_io_metrics::record_put_object_stage_duration( - rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_DISK_WAIT, - duration_ms, - ); - let position = if result.is_ok() { - let rank = successful_rename_completion_rank - .as_ref() - .map(|rank| rank.fetch_add(1, Ordering::Relaxed) + 1) - .unwrap_or(1); - if rank <= write_quorum { - rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_QUORUM_FIRST - } else { - rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_QUORUM_TAIL - } - } else { - rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_ERROR + let Some(disk) = disk else { + return Err(DiskError::DiskNotFound); }; - rustfs_io_metrics::record_put_rename_disk_wait_completion(position, duration_ms); - } - result - }) - .catch_unwind() + + let is_delete_marker = file_info.is_canonical_delete_marker(); + let mut local_file_info; + let file_info = if file_info.erasure.index == 0 { + local_file_info = file_info.clone(); + local_file_info.erasure.index = i + 1; + &local_file_info + } else { + file_info + }; + if file_info.erasure.index == 0 || (!is_delete_marker && !file_info.has_valid_erasure_geometry()) { + return Err(DiskError::FileCorrupt); + } + + // Test-only awaitable pause point right before the disk rename. + // A no-op immediately-ready future in production. + Self::rename_fanout_barrier(&dst_object, i, rename_fanout_barrier_phase::RENAME).await; + + if let Some(err) = Self::rename_injected_error(&dst_object, i) { + return Err(err); + } + + if let Some(scope) = publication_scope.as_ref() + && !scope.can_commit() + { + let _ = scope.mark_indeterminate(); + return Err(DiskError::other( + "scanner publication commit scope deadline or cancellation reached", + )); + } + + let disk_wait_started = rustfs_io_metrics::put_stage_timer(); + dispatched = true; + let result = disk + .rename_data_borrowed_with_fence( + &src_bucket, + &src_object, + file_info, + &dst_bucket, + &dst_object, + scanner_publication_lease_token, + ) + .await; + #[cfg(test)] + if result.is_ok() { + rollback_fault_injection::after_rename(&dst_object, i)?; + } + if let Some(disk_wait_started) = disk_wait_started { + let duration_ms = disk_wait_started.elapsed().as_secs_f64() * 1000.0; + rustfs_io_metrics::record_put_object_stage_duration( + rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_DISK_WAIT, + duration_ms, + ); + let position = if result.is_ok() { + let rank = successful_rename_completion_rank + .as_ref() + .map(|rank| rank.fetch_add(1, Ordering::Relaxed) + 1) + .unwrap_or(1); + if rank <= write_quorum { + rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_QUORUM_FIRST + } else { + rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_QUORUM_TAIL + } + } else { + rustfs_io_metrics::PUT_RENAME_DISK_WAIT_COMPLETION_POSITION_ERROR + }; + rustfs_io_metrics::record_put_rename_disk_wait_completion(position, duration_ms); + } + result + }) + .catch_unwind() + .await; + (dispatched, result) + } }); let results = join_all(futures).await; + #[cfg(test)] + rollback_fault_injection::after_fanout(&fanout_dst_object); (results, fanout_file_infos) }); @@ -4920,12 +5016,18 @@ impl SetDisks { rustfs_io_metrics::PUT_STAGE_SET_DISK_RENAME_QUORUM_WAIT, quorum_wait_started, ); - let (results, mut file_infos) = fanout_result.map_err(|_| DiskError::Unexpected)?; + let (results, mut file_infos) = match fanout_result { + Ok(result) => result, + Err(_) => { + record_indeterminate_rename(disks, (&dst_bucket, &dst_object), rollback_receipt).await; + return Err(DiskError::Unexpected); + } + }; if rustfs_io_metrics::put_stage_metrics_enabled() { let mut fanout_success = 0; let mut fanout_error = 0; let mut fanout_panic = 0; - for result in &results { + for (_, result) in &results { match result { Ok(Ok(_)) => fanout_success += 1, Ok(Err(_)) => fanout_error += 1, @@ -4941,7 +5043,9 @@ impl SetDisks { ); } - for (idx, result) in results.iter().enumerate() { + let mut dispatched = Vec::with_capacity(results.len()); + for (idx, (was_dispatched, result)) in results.iter().enumerate() { + dispatched.push(*was_dispatched); match result { Ok(Ok(res)) => { data_dirs[idx] = res.rollback_data_dir.or(res.old_data_dir); @@ -5003,7 +5107,16 @@ impl SetDisks { ); } - rollback_failed_rename(disks, file_infos, &errs, &data_dirs, (&dst_bucket, &dst_object), rollback_receipt).await; + rollback_failed_rename( + disks, + file_infos, + &errs, + &dispatched, + &data_dirs, + (&dst_bucket, &dst_object), + rollback_receipt, + ) + .await; return Err(ret_err); } @@ -6961,6 +7074,9 @@ pub(in crate::set_disk) mod rollback_fault_injection { pub(in crate::set_disk) enum Fault { Io, Panic, + IoAfterRename, + PanicAfterRename, + CoordinatorPanic, } fn registry() -> &'static Mutex> { @@ -6998,6 +7114,30 @@ pub(in crate::set_disk) mod rollback_fault_injection { _ => Ok(()), } } + + pub(super) fn after_rename(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::IoAfterRename)) if target == disk_index => Err(DiskError::FaultyDisk), + Some((target, Fault::PanicAfterRename)) if target == disk_index => panic!("injected panic after rename mutation"), + _ => Ok(()), + } + } + + pub(super) fn after_fanout(object: &str) { + let fault = registry() + .lock() + .expect("rollback registry should not poison") + .get(object) + .copied(); + if matches!(fault, Some((_, Fault::CoordinatorPanic))) { + panic!("injected rename coordinator panic"); + } + } } /// Test-only per-disk call counters for the metadata fan-out (backlog#1325, @@ -10722,6 +10862,7 @@ mod tests { } Some(rollback_fault_injection::Fault::Panic) => RenameRollbackOutcome::Panicked, None => RenameRollbackOutcome::Succeeded, + Some(_) => unreachable!("matrix only injects undo faults"), } } else { RenameRollbackOutcome::Succeeded @@ -10736,7 +10877,7 @@ mod tests { .join(STORAGE_FORMAT_FILE_BACKUP); assert_eq!( backup.exists(), - outcome.outcome.failed(), + outcome.outcome.needs_recovery(), "failed undo must retain its only old-version backup" ); } @@ -10799,94 +10940,188 @@ mod tests { #[tokio::test] #[serial_test::serial(capacity_dirty_scope)] - async fn rename_rollback_incomplete_preserves_overwrite_data_dirs_and_staging() { + async fn rename_data_early_ack_post_mutation_tail_error_never_rolls_back_commit() { 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}"); + for fault in [ + rollback_fault_injection::Fault::IoAfterRename, + rollback_fault_injection::Fault::PanicAfterRename, + ] { + let bucket = "rename-tail-unknown"; + let object = format!("tail-{fault:?}"); 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()), - ) + let _fault = rollback_fault_injection::arm(&object, 0, fault); + let barrier = rename_fanout_barrier::arm(&object, 0, rename_fanout_barrier_phase::RENAME); + 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), + true, + 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!("tail barrier must precede quorum ACK"), + } + }) + .await + .expect("tail reaches the barrier"); + let commit = tokio::time::timeout(BARRIER_PAUSE_GUARD, rename) .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")) + .expect("quorum must ACK before tail release") + .expect("three disks commit"); + barrier.release(); + let tail = commit + .tail_drain + .expect("early ACK owns a tail") + .await + .expect("tail coordinator joins") + .expect("committed tail reports convergence"); + assert_eq!(tail.convergence, RenameConvergence::PartialCommit); + assert!(receipt.0.get().is_none(), "post-ACK errors must never start rollback"); + for dir in &dirs { + let reopened = reopen_local_disk(dir).await; + let stored = reopened + .read_version( + "", + bucket, + &object, + "", + &ReadOptions { + read_data: true, + ..Default::default() + }, + ) + .await + .expect("all actual writes survive despite a lost tail acknowledgement"); + assert_eq!(stored.data.as_deref(), Some(b"inline-body".as_slice())); + } + } + }) + .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 fault in [ + rollback_fault_injection::Fault::Io, + rollback_fault_injection::Fault::IoAfterRename, + rollback_fault_injection::Fault::PanicAfterRename, + rollback_fault_injection::Fault::CoordinatorPanic, + ] { + for early_ack in [false, true] { + let bucket = "rollback-data-dirs"; + let object = format!("object-{early_ack}-{fault:?}"); + 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 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 + .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()) - .join("part.1"); - assert_eq!( - tokio::fs::read(staged) - .await - .expect("failed-write staging must remain available"), - b"new-data" + .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, fault); + 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()); + if !matches!(fault, rollback_fault_injection::Fault::Io) { + assert!( + matches!( + receipt.0.get().expect("indeterminate report").disks[0].outcome, + RenameRollbackOutcome::Indeterminate(_) + ), + "post-mutation failure must not be classified as unattempted" + ); + } + let coordinator_failed = matches!(fault, rollback_fault_injection::Fault::CoordinatorPanic); + 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 || (coordinator_failed && idx == 1), + "unknown mutations must retain the old-version 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 || (coordinator_failed && idx == 1) { + new_data_dir + } else { + old_data_dir + }) ); } - 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 })); } } }) @@ -10896,34 +11131,78 @@ mod tests { #[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"), + for cancel_caller in [false, true] { + let bucket = "rename-rollback-barrier"; + let object = if cancel_caller { + "rollback-barrier-cancelled" + } else { + "rollback-barrier-object" + }; + let (dirs, disks) = call_counter_local_disks(bucket, 4).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + 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 metadata should be staged"); } - }) - .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"); + 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"); + if cancel_caller { + drop(rename); + barrier.release(); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + while receipt.0.get().is_none() { + tokio::task::yield_now().await; + } + }) + .await + .expect("cancelled caller must not cancel rollback accounting"); + } else { + barrier.release(); + assert!(rename.await.is_err()); + } + assert!(receipt.is_incomplete(), "drained undo failure must survive in the receipt"); + for dir in dirs.iter().skip(1) { + let reopened = reopen_local_disk(dir).await; + let restored = reopened + .read_version( + "", + bucket, + object, + "", + &ReadOptions { + read_data: true, + ..Default::default() + }, + ) + .await + .expect("old version must remain readable after caller cancellation"); + assert_eq!(restored.data.as_deref(), Some(b"old-inline-body".as_slice())); + } + } } #[tokio::test]