diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 3bbe8b595..d17c1c751 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -138,6 +138,7 @@ use tokio::task::JoinSet; pub(in crate::set_disk) const EVENT_SET_DISK_READ: &str = "set_disk_read"; pub(in crate::set_disk) const ENV_RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP: &str = "RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP"; +pub(in crate::set_disk) const ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE: &str = "RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE"; /// Default reader-setup strategy for the GET read path (rustfs/backlog#1215, /// #1159, #923). /// @@ -3181,6 +3182,7 @@ pub(in crate::set_disk) struct RenameDataCommit { pub(in crate::set_disk) cleanup_disks: Vec>, pub(in crate::set_disk) old_current_size: Option, pub(in crate::set_disk) committed_file_info: FileInfo, + pub(in crate::set_disk) tail_drain: Option>, } #[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")] @@ -3192,6 +3194,10 @@ type RenameDataLegacyTuple = ( Option, ); +fn put_rename_early_ack_enabled() -> bool { + rustfs_utils::get_env_bool(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, false) +} + impl RenameDataCommit { #[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")] fn into_legacy_tuple(self) -> RenameDataLegacyTuple { @@ -3290,6 +3296,49 @@ impl SetDisks { } } + fn rename_data_commit_from_observations( + disks: &[Option], + file_infos: &[FileInfo], + disk_versions: &[Option>], + errs: &[Option], + cleanup_data_dirs: &[Option], + old_current_sizes: &[Option], + write_quorum: usize, + ) -> disk::error::Result { + let data_dir = Self::reduce_common_data_dir(cleanup_data_dirs, write_quorum); + let convergence = Self::classify_rename_convergence(disk_versions, errs); + let old_current_size = Self::reduce_common_old_current_size(old_current_sizes, write_quorum); + let online_disks = Self::eval_disks(disks, errs); + let committed_slot = online_disks.iter().position(Option::is_some).ok_or(DiskError::Unexpected)?; + let committed_file_info = file_infos.get(committed_slot).cloned().ok_or(DiskError::Unexpected)?; + let cleanup_disks = if let Some(data_dir) = data_dir { + disks + .iter() + .zip(errs.iter()) + .zip(cleanup_data_dirs.iter()) + .map(|((disk, err), old_data_dir)| { + if err.is_none() && *old_data_dir == Some(data_dir) { + disk.clone() + } else { + None + } + }) + .collect() + } else { + vec![None; disks.len()] + }; + + Ok(RenameDataCommit { + online_disks, + convergence, + data_dir, + cleanup_disks, + old_current_size, + committed_file_info, + tail_drain: None, + }) + } + pub(in crate::set_disk) async fn abort_quota_reservation_after_fence( reservation: crate::bucket::quota::reservation::QuotaReservation, disks: &[Option], @@ -3327,6 +3376,274 @@ impl SetDisks { .map(RenameDataCommit::into_legacy_tuple) } + #[tracing::instrument(level = "debug", skip(disks, file_infos))] + async fn rename_data_owned_early_ack( + disks: &[Option], + src_bucket: &str, + src_object: &str, + file_infos: Vec, + dst_bucket: &str, + dst_object: &str, + write_quorum: usize, + ) -> disk::error::Result { + if let Some(file_info) = disks + .iter() + .zip(file_infos.iter()) + .find_map(|(disk, file_info)| disk.as_ref().map(|_| file_info)) + { + if file_info.is_canonical_delete_marker() { + file_info.validate_for_metadata_read()?; + } else { + file_info.validate_for_erasure_write()?; + } + } + + let disk_count = disks.len(); + let fanout_disks = disks.to_vec(); + let coordinator_disks = fanout_disks.clone(); + let src_bucket = Arc::new(src_bucket.to_string()); + let src_object = Arc::new(src_object.to_string()); + let dst_bucket = Arc::new(dst_bucket.to_string()); + let dst_object = Arc::new(dst_object.to_string()); + let (commit_tx, commit_rx) = tokio::sync::oneshot::channel(); + + let tail_drain = tokio::spawn({ + let fanout_src_bucket = src_bucket.clone(); + let fanout_src_object = src_object.clone(); + let fanout_dst_bucket = dst_bucket.clone(); + let fanout_dst_object = dst_object.clone(); + async move { + let successful_rename_completion_rank = + rustfs_io_metrics::put_stage_metrics_enabled().then(|| Arc::new(AtomicUsize::new(0))); + let mut tasks = JoinSet::new(); + for (i, (disk, file_info)) in fanout_disks.into_iter().zip(file_infos.iter()).enumerate() { + let src_bucket = fanout_src_bucket.clone(); + let src_object = fanout_src_object.clone(); + let dst_bucket = fanout_dst_bucket.clone(); + let dst_object = fanout_dst_object.clone(); + 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 { + #[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); + } + + 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); + } + + let disk_wait_started = rustfs_io_metrics::put_stage_timer(); + let result = disk + .rename_data_borrowed(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object) + .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 + }; + rustfs_io_metrics::record_put_rename_disk_wait_completion(position, duration_ms); + } + result + }) + .catch_unwind() + .await; + (i, result) + }); + } + + let mut commit_tx = Some(commit_tx); + let mut sent_commit = false; + let mut success_count = 0usize; + let mut fanout_panic = 0usize; + let mut results_seen = 0usize; + let mut errs = vec![Some(DiskError::DiskNotFound); 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]; + let mut old_current_sizes = vec![None; disk_count]; + + while let Some(joined) = tasks.join_next().await { + results_seen += 1; + match joined { + Ok((idx, Ok(Ok(res)))) => { + 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; + old_current_sizes[idx] = res.old_current_size; + errs[idx] = None; + success_count += 1; + } + Ok((idx, Ok(Err(err)))) => { + errs[idx] = Some(err); + } + Ok((idx, Err(_))) => { + errs[idx] = Some(DiskError::Unexpected); + fanout_panic += 1; + } + Err(_) => { + fanout_panic += 1; + } + } + + if !sent_commit && success_count >= write_quorum { + let snapshot_commit = Self::rename_data_commit_from_observations( + &coordinator_disks, + &file_infos, + &disk_versions, + &errs, + &cleanup_data_dirs, + &old_current_sizes, + write_quorum, + ); + if let Some(commit_tx) = commit_tx.take() { + let _ = commit_tx.send(snapshot_commit); + } + sent_commit = true; + } + } + + 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); + rustfs_io_metrics::record_put_rename_quorum_wait_fanout( + results_seen, + write_quorum, + fanout_success, + fanout_error, + fanout_panic, + ); + } + + 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; + if let Some(commit_tx) = commit_tx.take() { + let _ = commit_tx.send(Err(ret_err)); + } + return; + } + + let mut backup_reclaims = Vec::new(); + for (idx, disk) in coordinator_disks.iter().enumerate() { + if errs[idx].is_some() { + continue; + } + let Some(rollback_dir) = data_dirs[idx] else { + continue; + }; + if cleanup_data_dirs[idx] == Some(rollback_dir) { + continue; + } + let Some(disk) = disk.clone() else { + continue; + }; + let dst_bucket = fanout_dst_bucket.clone(); + let dst_object = fanout_dst_object.clone(); + backup_reclaims.push(tokio::spawn(async move { + let backup_path = format!("{dst_object}/{rollback_dir}/{STORAGE_FORMAT_FILE_BACKUP}"); + disk.delete(&dst_bucket, &backup_path, DeleteOptions::default()).await + })); + } + for result in join_all(backup_reclaims).await { + match result { + Ok(Ok(())) => {} + Ok(Err(DiskError::FileNotFound | DiskError::VolumeNotFound)) => {} + Ok(Err(err)) => { + warn!( + dst_bucket = %fanout_dst_bucket, + dst_object = %fanout_dst_object, + error = %err, + "rollback backup reclamation failed after committed rename" + ); + } + Err(join_err) => { + warn!( + dst_bucket = %fanout_dst_bucket, + dst_object = %fanout_dst_object, + error = %join_err, + "rollback backup reclamation task failed after committed rename" + ); + } + } + } + } + }); + + let quorum_wait_started = rustfs_io_metrics::put_stage_timer(); + let commit = commit_rx.await.map_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, + ); + commit.map(|mut commit| { + commit.tail_drain = Some(tail_drain); + commit + }) + } + #[tracing::instrument(level = "debug", skip(disks, file_infos))] pub(in crate::set_disk) async fn rename_data_owned( disks: &[Option], @@ -3337,6 +3654,18 @@ impl SetDisks { dst_object: &str, write_quorum: usize, ) -> disk::error::Result { + if put_rename_early_ack_enabled() { + return Self::rename_data_owned_early_ack( + disks, + src_bucket, + src_object, + file_infos, + dst_bucket, + dst_object, + write_quorum, + ) + .await; + } if let Some(file_info) = disks .iter() .zip(file_infos.iter()) @@ -3411,6 +3740,10 @@ impl SetDisks { // 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); + } + let disk_wait_started = rustfs_io_metrics::put_stage_timer(); let result = disk .rename_data_borrowed(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object) @@ -3679,6 +4012,7 @@ impl SetDisks { cleanup_disks, old_current_size, committed_file_info, + tail_drain: None, }) } @@ -3890,6 +4224,20 @@ impl SetDisks { None } + /// Test-only fault-injection seam for the rename-data fan-out. In production + /// this is inlined to `None`; only the `#[cfg(test)]` variant consults the + /// fault registry after the per-disk barrier and before mutating the disk. + #[cfg(test)] + fn rename_injected_error(object: &str, disk_index: usize) -> Option { + rename_fault_injection::injected_error(object, disk_index) + } + + #[cfg(not(test))] + #[inline(always)] + fn rename_injected_error(_object: &str, _disk_index: usize) -> Option { + None + } + /// Test-only seam that records one per-disk `read_version` metadata RPC for /// the call-counter registry (backlog#1325). In production this is inlined to /// nothing and adds no behavior; only the `#[cfg(test)]` variant touches the @@ -5314,6 +5662,60 @@ pub(in crate::set_disk) mod cleanup_fault_injection { } } +/// Test-only rename fault-injection seam for the set-disk rename fan-out. +/// +/// This entire module is `#[cfg(test)]`, so it never compiles into production. +/// The registry is keyed by object name because the disk mutations happen in +/// spawned fan-out tasks; each test uses a unique object and clears the +/// registry through the returned guard. +#[cfg(test)] +pub(in crate::set_disk) mod rename_fault_injection { + use super::DiskError; + use std::sync::{Mutex, OnceLock}; + + #[derive(Default)] + struct FaultState { + object: Option, + fail_indices: Vec, + } + + fn state() -> &'static Mutex { + static STATE: OnceLock> = OnceLock::new(); + STATE.get_or_init(|| Mutex::new(FaultState::default())) + } + + /// Guard that clears the fault registry on drop. + pub struct FaultGuard; + + impl Drop for FaultGuard { + fn drop(&mut self) { + if let Ok(mut s) = state().lock() { + *s = FaultState::default(); + } + } + } + + /// Force rename-data mutations for `object` on the given disk indices to + /// fail with a transient error, until the returned guard is dropped. + #[must_use] + pub fn fail_rename_on(object: &str, indices: &[usize]) -> FaultGuard { + let mut s = state().lock().expect("rename fault registry poisoned"); + s.object = Some(object.to_string()); + s.fail_indices = indices.to_vec(); + FaultGuard + } + + pub(super) fn injected_error(object: &str, disk_index: usize) -> Option { + let s = state().lock().expect("rename fault registry poisoned"); + match &s.object { + Some(target) if target == object && s.fail_indices.contains(&disk_index) => { + Some(DiskError::other("injected rename failure (test-only)")) + } + _ => None, + } + } +} + /// Test-only per-disk call counters for the metadata fan-out (backlog#1325, /// serving the RPC-count assertions of #1309 / #1314 / #1315). /// @@ -7910,6 +8312,293 @@ mod tests { drop(dirs); } + #[tokio::test] + #[serial_test::serial(rename_quorum_ack)] + async fn rename_data_early_ack_returns_after_write_quorum_and_drains_tail_success() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + const DISKS: usize = 4; + let bucket = "rename-early-ack-tail-success-bucket"; + let object = "rename-early-ack-tail-success-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + let file_infos = rename_commit_fileinfos(object, DISKS, "early-tail-success-etag"); + let tracker = rename_fanout_barrier::observe_tasks(object); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + + let mut rename = Box::pin(SetDisks::rename_data( + &disks, + RUSTFS_META_TMP_BUCKET, + "source", + &file_infos, + bucket, + object, + 3, + )); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + tokio::select! { + () = barrier.wait_until_paused() => {} + result = rename.as_mut() => panic!("rename_data returned before the armed fan-out barrier: {result:?}"), + } + }) + .await + .expect("paused disk must reach the armed rename barrier"); + + rename + .await + .expect("early ACK must return once the other three disks satisfy write quorum"); + assert!( + tracker.running() >= 1, + "the paused tail disk must continue in the background after early ACK" + ); + + barrier.release(); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + while tracker.running() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("background tail disk must drain after release"); + + for (idx, dir) in dirs.iter().enumerate() { + let reopened = reopen_local_disk(dir).await; + let stored = reopened + .read_version("", bucket, object, "", &ReadOptions::default()) + .await + .unwrap_or_else(|err| panic!("disk {idx} must contain the early-ACK commit after reopen: {err:?}")); + assert_eq!( + stored.metadata.get("etag").map(String::as_str), + Some("early-tail-success-etag"), + "disk {idx} must expose the same committed metadata after background tail success and reopen" + ); + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial(rename_quorum_ack)] + async fn rename_data_early_ack_tail_failure_does_not_expose_partial_fresh_after_reopen() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + const DISKS: usize = 4; + let bucket = "rename-early-ack-tail-failure-bucket"; + let object = "rename-early-ack-tail-failure-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + let file_infos = rename_commit_fileinfos(object, DISKS, "early-tail-failure-etag"); + let tracker = rename_fanout_barrier::observe_tasks(object); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + let _fault = rename_fault_injection::fail_rename_on(object, &[0]); + + let mut rename = Box::pin(SetDisks::rename_data( + &disks, + RUSTFS_META_TMP_BUCKET, + "source", + &file_infos, + bucket, + object, + 3, + )); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + tokio::select! { + () = barrier.wait_until_paused() => {} + result = rename.as_mut() => panic!("rename_data returned before the armed fan-out barrier: {result:?}"), + } + }) + .await + .expect("paused disk must reach the armed rename barrier"); + + rename + .await + .expect("early ACK must return after write quorum before the tail failure"); + barrier.release(); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + while tracker.running() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("background tail failure must drain after release"); + + for (idx, dir) in dirs.iter().enumerate() { + let reopened = reopen_local_disk(dir).await; + let read = reopened.read_version("", bucket, object, "", &ReadOptions::default()).await; + if idx == 0 { + assert!( + matches!(read, Err(DiskError::FileNotFound | DiskError::FileVersionNotFound)), + "failed background tail disk must not expose a partial fresh commit after reopen: {read:?}" + ); + } else { + let stored = read.unwrap_or_else(|err| { + panic!("quorum disk {idx} must contain the early-ACK commit after reopen: {err:?}") + }); + assert_eq!( + stored.metadata.get("etag").map(String::as_str), + Some("early-tail-failure-etag"), + "quorum disk {idx} must expose the committed metadata after reopen" + ); + } + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial(rename_quorum_ack)] + async fn rename_data_early_ack_tail_failure_preserves_overwrite_after_reopen() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + const DISKS: usize = 4; + let bucket = "rename-early-ack-overwrite-tail-failure-bucket"; + let object = "rename-early-ack-overwrite-tail-failure-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).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(), "early-old-etag".to_string()); + for disk in disks.iter().flatten() { + disk.write_metadata(bucket, bucket, object, old.clone()) + .await + .expect("old metadata should be written before early overwrite"); + } + + let file_infos = rename_commit_fileinfos(object, DISKS, "early-new-etag"); + let tracker = rename_fanout_barrier::observe_tasks(object); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + let _fault = rename_fault_injection::fail_rename_on(object, &[0]); + + let mut rename = Box::pin(SetDisks::rename_data( + &disks, + RUSTFS_META_TMP_BUCKET, + "source", + &file_infos, + bucket, + object, + 3, + )); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + tokio::select! { + () = barrier.wait_until_paused() => {} + result = rename.as_mut() => panic!("rename_data returned before the armed fan-out barrier: {result:?}"), + } + }) + .await + .expect("paused disk must reach the armed rename barrier"); + + rename + .await + .expect("early ACK must return after write quorum before overwrite tail failure"); + barrier.release(); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + while tracker.running() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("background overwrite tail failure must drain after release"); + + for (idx, dir) in dirs.iter().enumerate() { + let reopened = reopen_local_disk(dir).await; + let stored = reopened + .read_version("", bucket, object, "", &ReadOptions::default()) + .await + .unwrap_or_else(|err| panic!("disk {idx} must have a readable version after early overwrite: {err:?}")); + let expected_etag = if idx == 0 { "early-old-etag" } else { "early-new-etag" }; + assert_eq!( + stored.metadata.get("etag").map(String::as_str), + Some(expected_etag), + "disk {idx} must keep the correct early-ACK overwrite visibility after reopen" + ); + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial(rename_quorum_ack)] + async fn rename_data_early_ack_background_drains_after_caller_cancellation() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + const DISKS: usize = 4; + let bucket = "rename-early-ack-cancel-bucket"; + let object = "rename-early-ack-cancel-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + let file_infos = rename_commit_fileinfos(object, DISKS, "early-cancel-etag"); + let tracker = rename_fanout_barrier::observe_tasks(object); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + + let mut rename = Box::pin(SetDisks::rename_data( + &disks, + RUSTFS_META_TMP_BUCKET, + "source", + &file_infos, + bucket, + object, + 3, + )); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + tokio::select! { + () = barrier.wait_until_paused() => {} + result = rename.as_mut() => panic!("rename_data returned before the armed fan-out barrier: {result:?}"), + } + }) + .await + .expect("paused disk must reach the armed rename barrier"); + + drop(rename); + barrier.release(); + tokio::time::timeout(BARRIER_PAUSE_GUARD, async { + while tracker.running() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("background rename fan-out must drain after caller cancellation"); + + for (idx, dir) in dirs.iter().enumerate() { + let reopened = reopen_local_disk(dir).await; + let stored = reopened + .read_version("", bucket, object, "", &ReadOptions::default()) + .await + .unwrap_or_else(|err| panic!("disk {idx} must contain the cancelled early-ACK commit after drain: {err:?}")); + assert_eq!( + stored.metadata.get("etag").map(String::as_str), + Some("early-cancel-etag"), + "disk {idx} must expose the background-drained commit after caller cancellation" + ); + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial(rename_quorum_ack)] + async fn rename_data_early_ack_strict_quorum_failure_rolls_back_fresh_after_reopen() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + const DISKS: usize = 4; + let bucket = "rename-early-ack-strict-rollback-bucket"; + let object = "rename-early-ack-strict-rollback-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + prepare_rename_source_dirs(&dirs, &disks, "source").await; + let file_infos = rename_commit_fileinfos(object, DISKS, "early-strict-rollback-etag"); + let _fault = rename_fault_injection::fail_rename_on(object, &[0]); + + SetDisks::rename_data(&disks, RUSTFS_META_TMP_BUCKET, "source", &file_infos, bucket, object, 4) + .await + .expect_err("three successful disks must fail an early-ACK strict write quorum of four"); + + for (idx, dir) in dirs.iter().enumerate() { + let reopened = reopen_local_disk(dir).await; + let read = reopened.read_version("", bucket, object, "", &ReadOptions::default()).await; + assert!( + matches!(read, Err(DiskError::FileNotFound | DiskError::FileVersionNotFound)), + "disk {idx} must not expose a fresh object after early-ACK strict rollback and reopen: {read:?}" + ); + } + }) + .await; + } + #[tokio::test] #[serial_test::serial(rename_quorum_ack)] async fn rename_data_commits_fresh_object_when_tail_disk_fails_after_write_quorum() { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index d4c535461..ca05840db 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -272,6 +272,7 @@ const MULTIPART_WRITE_QUORUM_RENAME_PART: &str = "rename_part"; const EVENT_SET_DISK_WRITE: &str = "set_disk_write"; const EVENT_SET_DISK_HEAL: &str = "set_disk_heal"; const EVENT_SET_DISK_COMMIT_TAIL_SLOW: &str = "set_disk_commit_tail_slow"; +const EVENT_SET_DISK_RENAME_TAIL_DRAIN_FAILED: &str = "set_disk_rename_tail_drain_failed"; const EVENT_SET_DISK_PUT_OBJECT_STAGE_SUMMARY: &str = "set_disk_put_object_stage_summary"; const SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS: u128 = 5_000; const ENV_RUSTFS_PUT_LARGE_BATCH_MIN_SIZE_BYTES: &str = "RUSTFS_PUT_LARGE_BATCH_MIN_SIZE_BYTES"; diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 35be8062c..0f418fa7f 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -2433,8 +2433,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let commit_object_lock_guard = object_lock_guard.take(); let detach_commit_owner = commit_object_lock_guard.is_some() || upload_guard.is_some() || quota_mutation_fence; let commit = async move { - let _object_lock_guard = commit_object_lock_guard; - let _upload_guard = upload_guard; + let mut _object_lock_guard = commit_object_lock_guard; + let mut _upload_guard = upload_guard; let mut quota_reservation = quota_reservation; let complete_tail_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); @@ -2570,6 +2570,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let op_old_dir = rename_commit.data_dir; let cleanup_disks = rename_commit.cleanup_disks; let committed_file_info = rename_commit.committed_file_info; + let rename_tail_drain = rename_commit.tail_drain; // Detach admission before any post-commit await: client cancellation // must not couple durable convergence repair to cleanup work. @@ -2628,7 +2629,30 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .invalidate_get_object_metadata_cache(&commit_bucket, &commit_object) .await; - drop(_object_lock_guard); // release the object lock before multipart cleanup tail IO. + if let Some(rename_tail_drain) = rename_tail_drain { + let object_lock_guard = _object_lock_guard.take(); + let upload_guard = _upload_guard.take(); + let tail_bucket = commit_bucket.clone(); + let tail_object = commit_object.clone(); + tokio::spawn(async move { + let _object_lock_guard = object_lock_guard; + let _upload_guard = upload_guard; + if let Err(err) = rename_tail_drain.await { + warn!( + event = EVENT_SET_DISK_RENAME_TAIL_DRAIN_FAILED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + state = "failed", + bucket = %tail_bucket, + object = %tail_object, + error = %err, + "rename tail drain failed" + ); + } + }); + } else { + drop(_object_lock_guard.take()); // release the object lock before multipart cleanup tail IO. + } #[cfg(test)] pause_multipart_commit(&commit_bucket, &commit_object, MultipartCommitPause::AfterObjectPublication).await; @@ -2685,7 +2709,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { ); } - drop(_upload_guard); + drop(_upload_guard.take()); Ok(ObjectInfo::from_file_info(&fi, &commit_bucket, &commit_object, commit_is_versioned)) }; diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 5430b5d1d..20ff815a7 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -2916,6 +2916,7 @@ impl SetDisks { let cleanup_disks = rename_commit.cleanup_disks; let old_current_size = rename_commit.old_current_size; let mut fi = rename_commit.committed_file_info; + let rename_tail_drain = rename_commit.tail_drain; // Do this before any post-commit await so request cancellation cannot // bypass best-effort admission. A process crash before admission // remains subject to the existing scanner reconciliation path. @@ -2954,11 +2955,38 @@ impl SetDisks { .invalidate_get_object_metadata_cache(&commit_bucket, &commit_object) .await; - // `rename_data` has completed the authoritative quorum commit. The - // exact old-data-dir reclamation below is best-effort space cleanup; - // it must not serialize the next operation on this object. - drop(_object_lock_guard); - drop(_bucket_lifecycle_guard); + // `rename_data` has completed the authoritative quorum commit. With + // the default-off early-ACK experiment, tail disk rename tasks may + // still be draining after quorum. Keep the namespace guards alive + // until that drain completes so the next same-object mutation cannot + // race a background tail rename. + if let Some(rename_tail_drain) = rename_tail_drain { + let object_lock_guard = _object_lock_guard; + let bucket_lifecycle_guard = _bucket_lifecycle_guard; + let tail_bucket = commit_bucket.clone(); + let tail_object = commit_object.clone(); + tokio::spawn(async move { + let _object_lock_guard = object_lock_guard; + let _bucket_lifecycle_guard = bucket_lifecycle_guard; + if let Err(err) = rename_tail_drain.await { + warn!( + event = EVENT_SET_DISK_RENAME_TAIL_DRAIN_FAILED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + state = "failed", + bucket = %tail_bucket, + object = %tail_object, + error = %err, + "rename tail drain failed" + ); + } + }); + } else { + // The exact old-data-dir reclamation below is best-effort space + // cleanup; it must not serialize the next operation on this object. + drop(_object_lock_guard); + drop(_bucket_lifecycle_guard); + } rustfs_io_metrics::record_put_object_stage_duration("set_disk_rename", duration_millis_f64(rename_stage_elapsed)); if (rename_stage_ms as u128) >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS { @@ -12700,7 +12728,7 @@ mod put_object_tmp_cleanup_tests { use super::hermetic_set_disks_support::hermetic_set_disks_isolated as hermetic_set_disks; use super::*; use crate::disk::DiskAPI as _; - use crate::set_disk::core::io_primitives::rename_fanout_barrier; + use crate::set_disk::core::io_primitives::{ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, rename_fanout_barrier}; use std::time::Duration; use tempfile::TempDir; use tokio::io::AsyncReadExt; @@ -13047,6 +13075,78 @@ mod put_object_tmp_cleanup_tests { assert_eq!(body, vec![b'2'; TEST_OBJECT_SIZE]); } + #[tokio::test] + #[serial_test::serial(rename_quorum_ack)] + async fn early_ack_tail_drain_retains_namespace_lock_until_background_rename_finishes() { + temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "put-early-ack-tail-lock"; + let object = "early-ack-tail-lock-object"; + for disk in &disk_stores { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + + let rename_tasks = rename_fanout_barrier::observe_tasks(object); + let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + let first_store = Arc::clone(&set_disks); + let first = tokio::spawn(async move { + let mut reader = PutObjReader::from_vec(vec![b'1'; TEST_OBJECT_SIZE]); + first_store + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + }); + tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused()) + .await + .expect("first PUT should pause one tail disk during rename"); + first + .await + .expect("first early-ACK PUT task should join before tail release") + .expect("first early-ACK PUT should return after write quorum"); + assert!( + rename_tasks.running() >= 1, + "tail rename must still be draining after the first PUT returns" + ); + + let second_namespace_barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::BeforeNamespace); + let second_store = Arc::clone(&set_disks); + let second = tokio::spawn(async move { + let mut reader = PutObjReader::from_vec(vec![b'2'; TEST_OBJECT_SIZE]); + second_store + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + }); + second_namespace_barrier.release_and_wait_until_namespace_pending().await; + tokio::task::yield_now().await; + assert!( + !second.is_finished(), + "second writer must remain blocked until the first early-ACK tail drain releases the namespace lock" + ); + + rename_barrier.release(); + drop(rename_barrier); + tokio::time::timeout(Duration::from_secs(30), async { + while rename_tasks.running() != 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("early-ACK tail rename should drain after release"); + second + .await + .expect("second overwrite task should join") + .expect("second overwrite should commit after the tail drain releases the namespace lock"); + + let mut reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("the latest overwrite should be readable"); + let mut body = Vec::new(); + reader.stream.read_to_end(&mut body).await.expect("latest body should drain"); + assert_eq!(body, vec![b'2'; TEST_OBJECT_SIZE]); + }) + .await; + } + #[tokio::test] async fn put_object_no_lock_aborts_after_outer_namespace_lock_loss() { let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;