fix(ecstore): retain indeterminate rename recovery evidence

This commit is contained in:
overtrue
2026-09-05 09:31:56 +08:00
parent dbe0073362
commit 19690eba7b
+522 -243
View File
@@ -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<DiskStore>],
mut file_infos: Vec<FileInfo>,
file_infos: Vec<FileInfo>,
errs: &[Option<DiskError>],
dispatched: &[bool],
rollback_dirs: &[Option<Uuid>],
dst: (&str, &str),
receipt: Option<RenameRollbackReceipt>,
) {
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<DiskStore>], dst: (&str, &str), receipt: Option<RenameRollbackReceipt>) {
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<DiskStore>],
outcomes: Vec<RenameRollbackDiskOutcome>,
dst: (&str, &str),
receipt: Option<RenameRollbackReceipt>,
) {
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<HashMap<String, (usize, Fault)>> {
@@ -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]