mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-29 08:27:06 +00:00
fix(ecstore): preserve rename data rollback metadata (#4107)
This commit is contained in:
@@ -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)
|
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<Option<String>> = 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) {
|
fn log_startup_disk_io_error(stage: &str, path: &Path, err: &IoError) {
|
||||||
warn!(
|
warn!(
|
||||||
event = EVENT_DISK_LOCAL_STARTUP_CLEANUP,
|
event = EVENT_DISK_LOCAL_STARTUP_CLEANUP,
|
||||||
@@ -3142,6 +3170,46 @@ impl DiskAPI for LocalDisk {
|
|||||||
return Err(err);
|
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 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() {
|
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;
|
let _ = self.delete_file(&dst_volume_dir, dst_data_path, false, false).await;
|
||||||
@@ -3159,29 +3227,6 @@ impl DiskAPI for LocalDisk {
|
|||||||
return Err(err);
|
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 let Some(src_file_path_parent) = src_file_path.parent() {
|
||||||
if src_volume != super::RUSTFS_META_MULTIPART_BUCKET {
|
if src_volume != super::RUSTFS_META_MULTIPART_BUCKET {
|
||||||
let _ = remove_std(src_file_path_parent);
|
let _ = remove_std(src_file_path_parent);
|
||||||
@@ -3240,6 +3285,17 @@ impl DiskAPI for LocalDisk {
|
|||||||
.truncate(true)
|
.truncate(true)
|
||||||
.open(&src)?;
|
.open(&src)?;
|
||||||
std::io::Write::write_all(&mut f, &new_buf)?;
|
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) {
|
match std::fs::rename(&src, &dst) {
|
||||||
Ok(()) => Ok(()),
|
Ok(()) => Ok(()),
|
||||||
Err(err) if err.kind() == std::io::ErrorKind::NotFound && !src.exists() => 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)),
|
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<uuid::Uuid>, Option<Bytes>), std::io::Error>((old_data_dir, has_dst_buf))
|
Ok::<(Option<uuid::Uuid>, Option<Bytes>), std::io::Error>((old_data_dir, has_dst_buf))
|
||||||
})
|
})
|
||||||
.await
|
.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
|
if let Some(old_data_dir) = opts.old_data_dir
|
||||||
&& opts.undo_write
|
&& opts.undo_write
|
||||||
{
|
{
|
||||||
@@ -3643,13 +3680,17 @@ impl DiskAPI for LocalDisk {
|
|||||||
file_path.as_path(),
|
file_path.as_path(),
|
||||||
Path::new(format!("{old_data_dir}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE_BACKUP}").as_str()),
|
Path::new(format!("{old_data_dir}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE_BACKUP}").as_str()),
|
||||||
]);
|
]);
|
||||||
let dst_path = path_join(&[
|
let dst_path = path_join(&[file_path.as_path(), Path::new(STORAGE_FORMAT_FILE)]);
|
||||||
file_path.as_path(),
|
|
||||||
Path::new(format!("{path}{SLASH_SEPARATOR}{STORAGE_FORMAT_FILE}").as_str()),
|
|
||||||
]);
|
|
||||||
return rename_all(&src_path, &dst_path, file_path).await;
|
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
|
self.delete_file(&volume_dir, &xl_path, true, false).await
|
||||||
}
|
}
|
||||||
#[tracing::instrument(level = "debug", skip(self))]
|
#[tracing::instrument(level = "debug", skip(self))]
|
||||||
@@ -3836,6 +3877,338 @@ mod test {
|
|||||||
use std::task::{Context, Poll};
|
use std::task::{Context, Poll};
|
||||||
use tokio::io::{AsyncReadExt, ReadBuf};
|
use tokio::io::{AsyncReadExt, ReadBuf};
|
||||||
|
|
||||||
|
fn test_file_info(name: &str, version_id: Uuid, data_dir: Option<Uuid>, data: Option<Bytes>) -> 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<u8> {
|
||||||
|
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]
|
#[tokio::test]
|
||||||
async fn test_skip_access_checks() {
|
async fn test_skip_access_checks() {
|
||||||
// let arr = Vec::new();
|
// let arr = Vec::new();
|
||||||
|
|||||||
@@ -6280,12 +6280,16 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
use crate::disk::CHECK_PART_UNKNOWN;
|
use crate::disk::CHECK_PART_UNKNOWN;
|
||||||
use crate::disk::CHECK_PART_VOLUME_NOT_FOUND;
|
use crate::disk::CHECK_PART_VOLUME_NOT_FOUND;
|
||||||
|
use crate::disk::DiskOption;
|
||||||
use crate::disk::RUSTFS_META_BUCKET;
|
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;
|
||||||
|
use crate::disk::STORAGE_FORMAT_FILE_BACKUP;
|
||||||
use crate::disk::WalkDirOptions;
|
use crate::disk::WalkDirOptions;
|
||||||
use crate::disk::endpoint::Endpoint;
|
use crate::disk::endpoint::Endpoint;
|
||||||
use crate::disk::error::DiskError;
|
use crate::disk::error::DiskError;
|
||||||
use crate::disk::health_state::RuntimeDriveHealthState;
|
use crate::disk::health_state::RuntimeDriveHealthState;
|
||||||
|
use crate::disk::new_disk;
|
||||||
use crate::layout::endpoints::SetupType;
|
use crate::layout::endpoints::SetupType;
|
||||||
use crate::object_api::ObjectInfo;
|
use crate::object_api::ObjectInfo;
|
||||||
use crate::storage_api_contracts::{
|
use crate::storage_api_contracts::{
|
||||||
@@ -6295,6 +6299,7 @@ mod tests {
|
|||||||
use crate::store::init_format::save_format_file;
|
use crate::store::init_format::save_format_file;
|
||||||
use crate::store::list_objects::ListPathOptions;
|
use crate::store::list_objects::ListPathOptions;
|
||||||
use rustfs_filemeta::ErasureInfo;
|
use rustfs_filemeta::ErasureInfo;
|
||||||
|
use rustfs_filemeta::FileMeta;
|
||||||
use rustfs_filemeta::MetaCacheEntry;
|
use rustfs_filemeta::MetaCacheEntry;
|
||||||
use rustfs_filemeta::ReplicationState;
|
use rustfs_filemeta::ReplicationState;
|
||||||
use rustfs_lock::client::local::LocalClient;
|
use rustfs_lock::client::local::LocalClient;
|
||||||
@@ -6303,6 +6308,7 @@ mod tests {
|
|||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use tempfile::TempDir;
|
use tempfile::TempDir;
|
||||||
use time::OffsetDateTime;
|
use time::OffsetDateTime;
|
||||||
|
use tokio::fs;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn complete_part_error_maps_confirmed_missing_to_invalid_part() {
|
fn complete_part_error_maps_confirmed_missing_to_invalid_part() {
|
||||||
@@ -6510,6 +6516,92 @@ mod tests {
|
|||||||
(dir, endpoint, disk)
|
(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]
|
#[test]
|
||||||
fn disk_health_entry_returns_cached_value_within_ttl() {
|
fn disk_health_entry_returns_cached_value_within_ttl() {
|
||||||
let entry = DiskHealthEntry {
|
let entry = DiskHealthEntry {
|
||||||
|
|||||||
@@ -132,26 +132,21 @@ impl SetDisks {
|
|||||||
let fi = file_infos[i].clone();
|
let fi = file_infos[i].clone();
|
||||||
let old_data_dir = data_dirs[i];
|
let old_data_dir = data_dirs[i];
|
||||||
let disk = disk.clone();
|
let disk = disk.clone();
|
||||||
let src_bucket = src_bucket.clone();
|
let dst_bucket = dst_bucket.clone();
|
||||||
let src_object = src_object.clone();
|
let dst_object = dst_object.clone();
|
||||||
futures.push(tokio::spawn(async move {
|
futures.push(tokio::spawn(async move {
|
||||||
let _ = disk
|
disk.delete_version(
|
||||||
.delete_version(
|
&dst_bucket,
|
||||||
&src_bucket,
|
&dst_object,
|
||||||
&src_object,
|
fi,
|
||||||
fi,
|
false,
|
||||||
false,
|
DeleteOptions {
|
||||||
DeleteOptions {
|
undo_write: true,
|
||||||
undo_write: true,
|
old_data_dir,
|
||||||
old_data_dir,
|
..Default::default()
|
||||||
..Default::default()
|
},
|
||||||
},
|
)
|
||||||
)
|
.await
|
||||||
.await
|
|
||||||
.map_err(|e| {
|
|
||||||
debug!("rename_data delete_version err {:?}", e);
|
|
||||||
e
|
|
||||||
});
|
|
||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -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);
|
return Err(ret_err);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user