From c67eea7bfaf0a13ac5e90727206654655f6f4544 Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 30 Jun 2026 15:45:27 +0800 Subject: [PATCH] fix(ecstore): preserve rename data rollback metadata (#4107) --- crates/ecstore/src/disk/local.rs | 465 ++++++++++++++++++++++++--- crates/ecstore/src/set_disk/mod.rs | 92 ++++++ crates/ecstore/src/set_disk/write.rs | 51 +-- 3 files changed, 542 insertions(+), 66 deletions(-) diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 3ed958307..7911a1ea4 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -107,6 +107,34 @@ fn get_direct_io_read_threshold() -> usize { rustfs_utils::get_env_usize(ENV_RUSTFS_OBJECT_DIRECT_IO_READ_THRESHOLD, DEFAULT_RUSTFS_OBJECT_DIRECT_IO_READ_THRESHOLD) } +#[cfg(test)] +static RENAME_DATA_FAIL_BEFORE_OLD_METADATA_BACKUP: std::sync::Mutex> = std::sync::Mutex::new(None); + +#[cfg(test)] +fn set_rename_data_fail_before_old_metadata_backup(dst_path: &str) { + *RENAME_DATA_FAIL_BEFORE_OLD_METADATA_BACKUP + .lock() + .expect("test failpoint lock should not be poisoned") = Some(dst_path.to_string()); +} + +#[cfg(test)] +fn should_fail_before_old_metadata_backup(dst_path: &str) -> bool { + let mut target = RENAME_DATA_FAIL_BEFORE_OLD_METADATA_BACKUP + .lock() + .expect("test failpoint lock should not be poisoned"); + if target.as_deref() == Some(dst_path) { + target.take(); + true + } else { + false + } +} + +#[cfg(not(test))] +fn should_fail_before_old_metadata_backup(_dst_path: &str) -> bool { + false +} + fn log_startup_disk_io_error(stage: &str, path: &Path, err: &IoError) { warn!( event = EVENT_DISK_LOCAL_STARTUP_CLEANUP, @@ -3142,6 +3170,46 @@ impl DiskAPI for LocalDisk { return Err(err); } + if should_fail_before_old_metadata_backup(dst_path) { + if let Some((_, dst_data_path)) = has_data_dir_path.as_ref() { + let _ = self.delete_file(&dst_volume_dir, dst_data_path, false, false).await; + } + info!( + event = EVENT_DISK_LOCAL_RENAME_REJECTED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_DISK_LOCAL, + reason = "test_fail_before_old_metadata_backup", + "Disk local rename flow failed before metadata commit" + ); + return Err(DiskError::Unexpected); + } + + if let Some(old_data_dir) = has_old_data_dir + && let Some(dst_buf) = has_dst_buf.as_ref() + && let Err(err) = self + .write_all_private( + dst_volume, + &format!("{}/{}/{}", &dst_path, &old_data_dir, STORAGE_FORMAT_FILE_BACKUP), + dst_buf.clone().into(), + true, + &skip_parent, + ) + .await + { + if let Some((_, dst_data_path)) = has_data_dir_path.as_ref() { + let _ = self.delete_file(&dst_volume_dir, dst_data_path, false, false).await; + } + info!( + event = EVENT_DISK_LOCAL_RENAME_REJECTED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_DISK_LOCAL, + reason = "write_old_metadata_backup_failed", + error = ?err, + "Disk local rename flow failed" + ); + return Err(err); + } + if let Err(err) = rename_all(&src_file_path, &dst_file_path, &skip_parent).await { if let Some((_, dst_data_path)) = has_data_dir_path.as_ref() { let _ = self.delete_file(&dst_volume_dir, dst_data_path, false, false).await; @@ -3159,29 +3227,6 @@ impl DiskAPI for LocalDisk { return Err(err); } - if let Some(old_data_dir) = has_old_data_dir - && let Some(dst_buf) = has_dst_buf - && let Err(err) = self - .write_all_private( - dst_volume, - &format!("{}/{}/{}", &dst_path, &old_data_dir.to_string(), STORAGE_FORMAT_FILE), - dst_buf.into(), - true, - &skip_parent, - ) - .await - { - info!( - event = EVENT_DISK_LOCAL_RENAME_REJECTED, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_DISK_LOCAL, - reason = "write_old_metadata_failed", - error = ?err, - "Disk local rename flow failed" - ); - return Err(err); - } - if let Some(src_file_path_parent) = src_file_path.parent() { if src_volume != super::RUSTFS_META_MULTIPART_BUCKET { let _ = remove_std(src_file_path_parent); @@ -3240,6 +3285,17 @@ impl DiskAPI for LocalDisk { .truncate(true) .open(&src)?; std::io::Write::write_all(&mut f, &new_buf)?; + if let Some(old_dir) = old_data_dir.as_ref() + && let Some(ref buf) = has_dst_buf + && let Some(dst_parent) = dst.parent() + { + let old_path = dst_parent.join(old_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP); + if let Some(old_parent) = old_path.parent() { + std::fs::create_dir_all(old_parent)?; + } + std::fs::write(&old_path, buf).map_err(to_file_error)?; + } + match std::fs::rename(&src, &dst) { Ok(()) => Ok(()), Err(err) if err.kind() == std::io::ErrorKind::NotFound && !src.exists() => Ok(()), @@ -3253,17 +3309,6 @@ impl DiskAPI for LocalDisk { Err(err) => Err(to_file_error(err)), }?; - if let Some(old_dir) = old_data_dir.as_ref() - && let Some(ref buf) = has_dst_buf - && let Some(dst_parent) = dst.parent() - { - let old_path = dst_parent.join(old_dir.to_string()).join(STORAGE_FORMAT_FILE); - if let Some(old_parent) = old_path.parent() { - std::fs::create_dir_all(old_parent)?; - } - std::fs::write(&old_path, buf).map_err(to_file_error)?; - } - Ok::<(Option, Option), std::io::Error>((old_data_dir, has_dst_buf)) }) .await @@ -3628,14 +3673,6 @@ impl DiskAPI for LocalDisk { } } - if !meta.versions.is_empty() { - let buf = meta.marshal_msg()?; - return self - .write_all_meta(volume, format!("{path}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE}").as_str(), &buf, true) - .await; - } - - // opts.undo_write && opts.old_data_dir.is_some_and(f) if let Some(old_data_dir) = opts.old_data_dir && opts.undo_write { @@ -3643,13 +3680,17 @@ impl DiskAPI for LocalDisk { file_path.as_path(), Path::new(format!("{old_data_dir}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE_BACKUP}").as_str()), ]); - let dst_path = path_join(&[ - file_path.as_path(), - Path::new(format!("{path}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE}").as_str()), - ]); + let dst_path = path_join(&[file_path.as_path(), Path::new(STORAGE_FORMAT_FILE)]); return rename_all(&src_path, &dst_path, file_path).await; } + if !meta.versions.is_empty() { + let buf = meta.marshal_msg()?; + return self + .write_all_meta(volume, format!("{path}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE}").as_str(), &buf, true) + .await; + } + self.delete_file(&volume_dir, &xl_path, true, false).await } #[tracing::instrument(level = "debug", skip(self))] @@ -3836,6 +3877,338 @@ mod test { use std::task::{Context, Poll}; use tokio::io::{AsyncReadExt, ReadBuf}; + fn test_file_info(name: &str, version_id: Uuid, data_dir: Option, data: Option) -> FileInfo { + let size = data + .as_ref() + .map(|data| i64::try_from(data.len()).expect("test data length should fit i64")) + .unwrap_or(1); + FileInfo { + name: name.to_string(), + version_id: Some(version_id), + data_dir, + data, + size, + mod_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + } + } + + fn test_meta(fi: FileInfo) -> Vec { + let mut meta = FileMeta::default(); + meta.add_version(fi).expect("test metadata should accept file info"); + meta.marshal_msg().expect("test metadata should encode") + } + + async fn ensure_test_volume(disk: &LocalDisk, volume: &str) { + match disk.make_volume(volume).await { + Ok(()) | Err(DiskError::VolumeExists) => {} + Err(err) => panic!("test volume should be available: {err:?}"), + } + } + + #[tokio::test] + async fn test_rename_data_writes_old_metadata_backup_before_non_inline_undo() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "dir/object"; + let tmp_object = "tmp-write"; + let version_id = Uuid::parse_str("11111111-1111-1111-1111-111111111111").expect("version id should parse"); + let old_data_dir = Uuid::parse_str("22222222-2222-2222-2222-222222222222").expect("old data dir should parse"); + let new_data_dir = Uuid::parse_str("33333333-3333-3333-3333-333333333333").expect("new data dir should parse"); + + ensure_test_volume(&disk, bucket).await; + ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + + let old_fi = test_file_info(object, version_id, Some(old_data_dir), None); + let dst_object_dir = dir.path().join(bucket).join("dir/object"); + fs::create_dir_all(dst_object_dir.join(old_data_dir.to_string())) + .await + .expect("old data dir should be created"); + fs::write(dst_object_dir.join(STORAGE_FORMAT_FILE), test_meta(old_fi)) + .await + .expect("old metadata should be written"); + + let tmp_data_dir = dir + .path() + .join(RUSTFS_META_TMP_BUCKET) + .join(tmp_object) + .join(new_data_dir.to_string()); + fs::create_dir_all(&tmp_data_dir) + .await + .expect("new tmp data dir should be created"); + fs::write(tmp_data_dir.join("part.1"), b"new-data") + .await + .expect("new tmp data should be written"); + + let new_fi = test_file_info(object, version_id, Some(new_data_dir), None); + let resp = disk + .rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) + .await + .expect("rename_data should commit"); + + assert_eq!(resp.old_data_dir, Some(old_data_dir)); + assert!( + dst_object_dir + .join(old_data_dir.to_string()) + .join(STORAGE_FORMAT_FILE_BACKUP) + .exists() + ); + assert!( + !dst_object_dir + .join(old_data_dir.to_string()) + .join(STORAGE_FORMAT_FILE) + .exists() + ); + } + + #[tokio::test] + async fn test_rename_data_writes_old_metadata_backup_for_inline_overwrite() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "inline-object"; + let tmp_object = "tmp-inline-write"; + let version_id = Uuid::parse_str("12121212-1212-1212-1212-121212121212").expect("version id should parse"); + let old_data_dir = Uuid::parse_str("34343434-3434-3434-3434-343434343434").expect("old data dir should parse"); + + ensure_test_volume(&disk, bucket).await; + ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + + let old_fi = test_file_info(object, version_id, Some(old_data_dir), None); + let dst_object_dir = dir.path().join(bucket).join(object); + fs::create_dir_all(dst_object_dir.join(old_data_dir.to_string())) + .await + .expect("old data dir should be created"); + fs::write(dst_object_dir.join(STORAGE_FORMAT_FILE), test_meta(old_fi)) + .await + .expect("old metadata should be written"); + + let tmp_object_dir = dir.path().join(RUSTFS_META_TMP_BUCKET).join(tmp_object); + fs::create_dir_all(&tmp_object_dir) + .await + .expect("tmp object dir should be created"); + + let new_fi = test_file_info(object, version_id, None, Some(Bytes::from_static(b"inline-new"))); + let resp = disk + .rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) + .await + .expect("inline rename_data should commit"); + + assert_eq!(resp.old_data_dir, Some(old_data_dir)); + assert!( + dst_object_dir + .join(old_data_dir.to_string()) + .join(STORAGE_FORMAT_FILE_BACKUP) + .exists() + ); + assert!( + !dst_object_dir + .join(old_data_dir.to_string()) + .join(STORAGE_FORMAT_FILE) + .exists() + ); + } + + #[tokio::test] + async fn test_delete_version_undo_restores_backup_to_object_root() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "dir/object"; + let version_id = Uuid::parse_str("44444444-4444-4444-4444-444444444444").expect("version id should parse"); + let old_data_dir = Uuid::parse_str("55555555-5555-5555-5555-555555555555").expect("old data dir should parse"); + let new_data_dir = Uuid::parse_str("66666666-6666-6666-6666-666666666666").expect("new data dir should parse"); + + ensure_test_volume(&disk, bucket).await; + + let object_dir = dir.path().join(bucket).join("dir/object"); + fs::create_dir_all(object_dir.join(old_data_dir.to_string())) + .await + .expect("old backup dir should be created"); + fs::create_dir_all(object_dir.join(new_data_dir.to_string())) + .await + .expect("new data dir should be created"); + + let old_fi = test_file_info(object, version_id, Some(old_data_dir), None); + let old_meta = test_meta(old_fi); + let new_fi = test_file_info(object, version_id, Some(new_data_dir), None); + fs::write( + object_dir.join(old_data_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP), + old_meta.clone(), + ) + .await + .expect("old metadata backup should be written"); + fs::write(object_dir.join(STORAGE_FORMAT_FILE), test_meta(new_fi.clone())) + .await + .expect("new metadata should be written"); + + disk.delete_version( + bucket, + object, + new_fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir: Some(old_data_dir), + ..Default::default() + }, + ) + .await + .expect("undo should restore old metadata"); + + let restored_meta = fs::read(object_dir.join(STORAGE_FORMAT_FILE)) + .await + .expect("restored metadata should be readable"); + assert_eq!(restored_meta, old_meta); + assert!(!object_dir.join("dir/object").join(STORAGE_FORMAT_FILE).exists()); + } + + #[tokio::test] + async fn test_delete_version_undo_restores_backup_when_other_versions_remain() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "dir/object"; + let version_id = Uuid::parse_str("44444444-4444-4444-4444-444444444444").expect("version id should parse"); + let other_version_id = Uuid::parse_str("77777777-7777-7777-7777-777777777777").expect("version id should parse"); + let old_data_dir = Uuid::parse_str("55555555-5555-5555-5555-555555555555").expect("old data dir should parse"); + let new_data_dir = Uuid::parse_str("66666666-6666-6666-6666-666666666666").expect("new data dir should parse"); + let other_data_dir = Uuid::parse_str("88888888-8888-8888-8888-888888888888").expect("other data dir should parse"); + + ensure_test_volume(&disk, bucket).await; + + let object_dir = dir.path().join(bucket).join("dir/object"); + fs::create_dir_all(object_dir.join(old_data_dir.to_string())) + .await + .expect("old backup dir should be created"); + fs::create_dir_all(object_dir.join(new_data_dir.to_string())) + .await + .expect("new data dir should be created"); + + let old_fi = test_file_info(object, version_id, Some(old_data_dir), None); + let other_fi = test_file_info(object, other_version_id, Some(other_data_dir), None); + let mut old_meta = FileMeta::default(); + old_meta + .add_version(old_fi) + .expect("old metadata should accept old file info"); + old_meta + .add_version(other_fi.clone()) + .expect("old metadata should accept other file info"); + let old_meta = old_meta.marshal_msg().expect("old metadata should encode"); + + let new_fi = test_file_info(object, version_id, Some(new_data_dir), None); + let mut new_meta = FileMeta::default(); + new_meta + .add_version(new_fi.clone()) + .expect("new metadata should accept new file info"); + new_meta + .add_version(other_fi) + .expect("new metadata should accept other file info"); + + fs::write( + object_dir.join(old_data_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP), + old_meta.clone(), + ) + .await + .expect("old metadata backup should be written"); + fs::write( + object_dir.join(STORAGE_FORMAT_FILE), + new_meta.marshal_msg().expect("new metadata should encode"), + ) + .await + .expect("new metadata should be written"); + + disk.delete_version( + bucket, + object, + new_fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir: Some(old_data_dir), + ..Default::default() + }, + ) + .await + .expect("undo should restore old metadata"); + + let restored_meta = fs::read(object_dir.join(STORAGE_FORMAT_FILE)) + .await + .expect("restored metadata should be readable"); + assert_eq!(restored_meta, old_meta); + } + + #[tokio::test] + async fn test_rename_data_failure_before_metadata_commit_preserves_old_metadata() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "failpoint-object"; + let tmp_object = "tmp-object"; + let version_id = Uuid::parse_str("aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa").expect("version id should parse"); + let old_data_dir = Uuid::parse_str("bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb").expect("old data dir should parse"); + let new_data_dir = Uuid::parse_str("cccccccc-cccc-cccc-cccc-cccccccccccc").expect("version id should parse"); + + ensure_test_volume(&disk, bucket).await; + ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + + let old_fi = test_file_info(object, version_id, Some(old_data_dir), None); + let old_meta = test_meta(old_fi); + let object_dir = dir.path().join(bucket).join(object); + fs::create_dir_all(object_dir.join(old_data_dir.to_string())) + .await + .expect("old data dir should be created"); + fs::write(object_dir.join(STORAGE_FORMAT_FILE), old_meta.clone()) + .await + .expect("old metadata should be written"); + + let tmp_data_dir = dir + .path() + .join(RUSTFS_META_TMP_BUCKET) + .join(tmp_object) + .join(new_data_dir.to_string()); + fs::create_dir_all(&tmp_data_dir) + .await + .expect("new tmp data dir should be created"); + fs::write(tmp_data_dir.join("part.1"), b"new-data") + .await + .expect("new tmp data should be written"); + + set_rename_data_fail_before_old_metadata_backup(object); + let new_fi = test_file_info(object, version_id, Some(new_data_dir), None); + let result = disk + .rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) + .await; + + assert!(result.is_err()); + let current_meta = fs::read(object_dir.join(STORAGE_FORMAT_FILE)) + .await + .expect("old metadata should still be readable"); + assert_eq!(current_meta, old_meta); + assert!(!object_dir.join("object").join(STORAGE_FORMAT_FILE).exists()); + } + #[tokio::test] async fn test_skip_access_checks() { // let arr = Vec::new(); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 54dd752b8..abb1c17b3 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -6280,12 +6280,16 @@ mod tests { use super::*; use crate::disk::CHECK_PART_UNKNOWN; use crate::disk::CHECK_PART_VOLUME_NOT_FOUND; + use crate::disk::DiskOption; use crate::disk::RUSTFS_META_BUCKET; + use crate::disk::RUSTFS_META_TMP_BUCKET; use crate::disk::STORAGE_FORMAT_FILE; + use crate::disk::STORAGE_FORMAT_FILE_BACKUP; use crate::disk::WalkDirOptions; use crate::disk::endpoint::Endpoint; use crate::disk::error::DiskError; use crate::disk::health_state::RuntimeDriveHealthState; + use crate::disk::new_disk; use crate::layout::endpoints::SetupType; use crate::object_api::ObjectInfo; use crate::storage_api_contracts::{ @@ -6295,6 +6299,7 @@ mod tests { use crate::store::init_format::save_format_file; use crate::store::list_objects::ListPathOptions; use rustfs_filemeta::ErasureInfo; + use rustfs_filemeta::FileMeta; use rustfs_filemeta::MetaCacheEntry; use rustfs_filemeta::ReplicationState; use rustfs_lock::client::local::LocalClient; @@ -6303,6 +6308,7 @@ mod tests { use std::collections::HashMap; use tempfile::TempDir; use time::OffsetDateTime; + use tokio::fs; #[test] fn complete_part_error_maps_confirmed_missing_to_invalid_part() { @@ -6510,6 +6516,92 @@ mod tests { (dir, endpoint, disk) } + #[tokio::test] + async fn test_rename_data_quorum_failure_rolls_back_destination_object() { + let dir = tempfile::tempdir().expect("tempdir should be created"); + let disk_root = dir.path().join("disk0"); + fs::create_dir_all(&disk_root).await.expect("disk root should be created"); + let endpoint = Endpoint::try_from(disk_root.to_str().expect("disk path should be utf8")).expect("endpoint should parse"); + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("disk should be created"); + + let bucket = "bucket"; + let object = "object"; + let tmp_object = "tmp-object"; + let version_id = Uuid::parse_str("77777777-7777-7777-7777-777777777777").expect("version id should parse"); + let old_data_dir = Uuid::parse_str("88888888-8888-8888-8888-888888888888").expect("old data dir should parse"); + let new_data_dir = Uuid::parse_str("99999999-9999-9999-9999-999999999999").expect("new data dir should parse"); + + match disk.make_volume(bucket).await { + Ok(()) | Err(DiskError::VolumeExists) => {} + Err(err) => panic!("bucket should be available: {err:?}"), + } + match disk.make_volume(RUSTFS_META_TMP_BUCKET).await { + Ok(()) | Err(DiskError::VolumeExists) => {} + Err(err) => panic!("tmp bucket should be available: {err:?}"), + } + + let object_dir = disk_root.join(bucket).join(object); + fs::create_dir_all(object_dir.join(old_data_dir.to_string())) + .await + .expect("old data dir should be created"); + let mut old_fi = FileInfo::new(&format!("{bucket}/{object}"), 1, 1); + old_fi.name = object.to_string(); + old_fi.version_id = Some(version_id); + old_fi.data_dir = Some(old_data_dir); + old_fi.size = 1; + old_fi.mod_time = Some(OffsetDateTime::now_utc()); + let mut old_meta = FileMeta::default(); + old_meta.add_version(old_fi).expect("old metadata should accept file info"); + let old_meta_buf = old_meta.marshal_msg().expect("old metadata should encode"); + fs::write(object_dir.join(STORAGE_FORMAT_FILE), old_meta_buf.clone()) + .await + .expect("old metadata should be written"); + + let tmp_data_dir = disk_root + .join(RUSTFS_META_TMP_BUCKET) + .join(tmp_object) + .join(new_data_dir.to_string()); + fs::create_dir_all(&tmp_data_dir) + .await + .expect("new tmp data dir should be created"); + fs::write(tmp_data_dir.join("part.1"), b"new") + .await + .expect("new tmp part should be written"); + + let mut new_fi = FileInfo::new(&format!("{bucket}/{object}"), 1, 1); + new_fi.name = object.to_string(); + new_fi.version_id = Some(version_id); + new_fi.data_dir = Some(new_data_dir); + new_fi.size = 1; + new_fi.mod_time = Some(OffsetDateTime::now_utc()); + + let disks = vec![Some(disk), None]; + let file_infos = vec![new_fi.clone(), new_fi]; + let result = SetDisks::rename_data(&disks, RUSTFS_META_TMP_BUCKET, tmp_object, &file_infos, bucket, object, 2).await; + + assert!(result.is_err()); + let restored_meta = fs::read(object_dir.join(STORAGE_FORMAT_FILE)) + .await + .expect("destination metadata should remain readable"); + assert_eq!(restored_meta, old_meta_buf); + assert!(!object_dir.join(object).join(STORAGE_FORMAT_FILE).exists()); + assert!(!object_dir.join(new_data_dir.to_string()).exists()); + assert!( + !object_dir + .join(old_data_dir.to_string()) + .join(STORAGE_FORMAT_FILE_BACKUP) + .exists() + ); + } + #[test] fn disk_health_entry_returns_cached_value_within_ttl() { let entry = DiskHealthEntry { diff --git a/crates/ecstore/src/set_disk/write.rs b/crates/ecstore/src/set_disk/write.rs index 4b60a20a4..03b833d43 100644 --- a/crates/ecstore/src/set_disk/write.rs +++ b/crates/ecstore/src/set_disk/write.rs @@ -132,26 +132,21 @@ impl SetDisks { let fi = file_infos[i].clone(); let old_data_dir = data_dirs[i]; let disk = disk.clone(); - let src_bucket = src_bucket.clone(); - let src_object = src_object.clone(); + let dst_bucket = dst_bucket.clone(); + let dst_object = dst_object.clone(); futures.push(tokio::spawn(async move { - let _ = disk - .delete_version( - &src_bucket, - &src_object, - fi, - false, - DeleteOptions { - undo_write: true, - old_data_dir, - ..Default::default() - }, - ) - .await - .map_err(|e| { - debug!("rename_data delete_version err {:?}", e); - e - }); + disk.delete_version( + &dst_bucket, + &dst_object, + fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir, + ..Default::default() + }, + ) + .await })); } } @@ -171,7 +166,23 @@ impl SetDisks { ); } - let _ = join_all(futures).await; + let undo_results = join_all(futures).await; + let undo_error_count = undo_results + .iter() + .filter(|result| match result { + Err(_) | Ok(Err(_)) => true, + Ok(Ok(_)) => false, + }) + .count(); + if undo_error_count > 0 { + warn!( + target: "rustfs_ecstore::set_disk", + dst_bucket = %dst_bucket, + dst_object = %dst_object, + undo_error_count, + "rename_data quorum rollback reported errors" + ); + } return Err(ret_err); }