diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index 133de9a0a..9feb1d7b7 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -152,4 +152,7 @@ mod namespace_lock_quorum_test; #[cfg(test)] mod admin_timeout_regression_test; +#[cfg(test)] +mod overwrite_cleanup_regression_test; + pub mod tls_gen; diff --git a/crates/e2e_test/src/overwrite_cleanup_regression_test.rs b/crates/e2e_test/src/overwrite_cleanup_regression_test.rs new file mode 100644 index 000000000..f41113b12 --- /dev/null +++ b/crates/e2e_test/src/overwrite_cleanup_regression_test.rs @@ -0,0 +1,112 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use aws_sdk_s3::primitives::ByteStream; +use serial_test::serial; +use std::path::{Path, PathBuf}; +use uuid::Uuid; + +use crate::common::{RustFSTestEnvironment, init_logging}; + +const TEST_BUCKET: &str = "overwrite-cleanup-regression"; +const TEST_OBJECT: &str = "large-object.bin"; +const PAYLOAD_SIZE: usize = 512 * 1024; + +#[tokio::test(flavor = "multi_thread")] +#[serial] +async fn unversioned_overwrite_removes_previous_physical_data_dir() -> Result<(), Box> { + init_logging(); + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + + let client = env.create_s3_client(); + client.create_bucket().bucket(TEST_BUCKET).send().await?; + + let first_payload = vec![b'a'; PAYLOAD_SIZE]; + client + .put_object() + .bucket(TEST_BUCKET) + .key(TEST_OBJECT) + .body(ByteStream::from(first_payload)) + .send() + .await?; + + let object_dir = Path::new(&env.temp_dir).join(TEST_BUCKET).join(TEST_OBJECT); + let first_data_dirs = physical_data_dirs(&object_dir)?; + assert_eq!( + first_data_dirs.len(), + 1, + "first non-inline put should create exactly one physical data directory: {first_data_dirs:?}" + ); + + let second_payload = vec![b'z'; PAYLOAD_SIZE]; + client + .put_object() + .bucket(TEST_BUCKET) + .key(TEST_OBJECT) + .body(ByteStream::from(second_payload.clone())) + .send() + .await?; + + let second_data_dirs = physical_data_dirs(&object_dir)?; + assert_eq!( + second_data_dirs.len(), + 1, + "unversioned overwrite should clean the previous physical data directory: {second_data_dirs:?}" + ); + assert_ne!( + first_data_dirs[0], second_data_dirs[0], + "overwrite should replace the old physical data directory with the new one" + ); + + let current = client + .get_object() + .bucket(TEST_BUCKET) + .key(TEST_OBJECT) + .send() + .await? + .body + .collect() + .await? + .into_bytes(); + assert_eq!( + current.as_ref(), + second_payload.as_slice(), + "object data should remain readable after old data directory cleanup" + ); + + env.stop_server(); + Ok(()) +} + +fn physical_data_dirs(object_dir: &Path) -> Result, Box> { + let mut data_dirs = Vec::new(); + for entry in std::fs::read_dir(object_dir)? { + let entry = entry?; + if !entry.file_type()?.is_dir() { + continue; + } + + let path: PathBuf = entry.path(); + let Some(name) = path.file_name().and_then(|name| name.to_str()) else { + continue; + }; + if let Ok(uuid) = Uuid::parse_str(name) { + data_dirs.push(uuid); + } + } + data_dirs.sort_unstable(); + Ok(data_dirs) +} diff --git a/crates/filemeta/src/filemeta/inline_data.rs b/crates/filemeta/src/filemeta/inline_data.rs index 92e774fc2..98d3bb28f 100644 --- a/crates/filemeta/src/filemeta/inline_data.rs +++ b/crates/filemeta/src/filemeta/inline_data.rs @@ -13,8 +13,33 @@ // limitations under the License. use super::*; +use rustfs_utils::http::{SUFFIX_TRANSITION_STATUS, get_bytes}; use std::collections::HashSet; +fn physical_data_dir(version: &FileMetaShallowVersion) -> Result> { + if version.header.version_type != VersionType::Object { + return Ok(None); + } + + let version = FileMetaVersion::try_from(version.meta.as_slice())?; + let Some(obj) = version.object else { + return Ok(None); + }; + + if obj.inlinedata() { + return Ok(None); + } + + if let Some(status) = get_bytes(&obj.meta_sys, SUFFIX_TRANSITION_STATUS) + && status.as_slice() == TRANSITION_COMPLETE.as_bytes() + && !is_restored_object_on_disk(&obj.meta_user) + { + return Ok(None); + } + + Ok(obj.data_dir.filter(|dir| !dir.is_nil())) +} + impl FileMeta { pub fn find_unshared_data_dir_for_version(&self, version_id: Option) -> Option { let vid = version_id.unwrap_or_default(); @@ -22,19 +47,20 @@ impl FileMeta { let mut target_selected = false; let mut other_data_dirs = HashSet::new(); - for version in self - .versions - .iter() - .filter(|v| v.header.version_type == VersionType::Object && v.header.uses_data_dir()) - { + for version in self.versions.iter().filter(|v| v.header.version_type == VersionType::Object) { let is_target_version = version.header.version_id.unwrap_or_default() == vid; + let dir = match physical_data_dir(version) { + Ok(dir) => dir, + Err(_) => return None, + }; + if is_target_version { if target_selected { continue; } target_selected = true; - target_data_dir = FileMetaVersion::decode_data_dir_from_meta(&version.meta).unwrap_or_default(); + target_data_dir = dir; if let Some(dir) = target_data_dir && other_data_dirs.contains(&dir) { @@ -43,7 +69,6 @@ impl FileMeta { continue; } - let dir = FileMetaVersion::decode_data_dir_from_meta(&version.meta).unwrap_or_default(); if let Some(dir) = dir { if target_data_dir == Some(dir) { return None; @@ -57,14 +82,23 @@ impl FileMeta { pub fn shard_data_dir_count(&self, vid: &Option, data_dir: &Option) -> usize { let vid = vid.unwrap_or_default(); - self.versions + let Some(data_dir) = data_dir else { + return 0; + }; + + let mut count = 0; + for version in self + .versions .iter() - .filter(|v| { - v.header.version_type == VersionType::Object && v.header.version_id != Some(vid) && v.header.uses_data_dir() - }) - .map(|v| FileMetaVersion::decode_data_dir_from_meta(&v.meta).unwrap_or_default()) - .filter(|v| v == data_dir) - .count() + .filter(|v| v.header.version_type == VersionType::Object && v.header.version_id != Some(vid)) + { + match physical_data_dir(version) { + Ok(Some(dir)) if dir == *data_dir => count += 1, + Ok(_) => {} + Err(_) => return 1, + } + } + count } pub fn get_data_dirs(&self) -> Result>> { @@ -81,6 +115,9 @@ impl FileMeta { /// Count shared data directories pub fn shared_data_dir_count(&self, version_id: Option, data_dir: Option) -> usize { let vid = version_id.unwrap_or_default(); + let Some(data_dir) = data_dir else { + return 0; + }; if self.data.entries().unwrap_or_default() > 0 && self.find_inline_data_for_version(version_id).unwrap_or_default().is_some() @@ -88,14 +125,19 @@ impl FileMeta { return 0; } - self.versions + let mut count = 0; + for version in self + .versions .iter() - .filter(|v| { - v.header.version_type == VersionType::Object && v.header.version_id != Some(vid) && v.header.uses_data_dir() - }) - .filter_map(|v| FileMetaVersion::decode_data_dir_from_meta(&v.meta).ok()) - .filter(|&dir| dir.is_some() && dir == data_dir) - .count() + .filter(|v| v.header.version_type == VersionType::Object && v.header.version_id != Some(vid)) + { + match physical_data_dir(version) { + Ok(Some(dir)) if dir == data_dir => count += 1, + Ok(_) => {} + Err(_) => return 1, + } + } + count } } @@ -106,6 +148,16 @@ mod tests { use time::format_description::well_known::Rfc3339; use time::{Duration, OffsetDateTime}; + fn make_plain_file_info(version_id: Uuid, data_dir: Uuid) -> FileInfo { + make_file_info_with_metadata(version_id, data_dir, HashMap::from([("etag".to_string(), format!("etag-{version_id}"))])) + } + + fn make_plain_null_version_file_info(data_dir: Uuid) -> FileInfo { + let mut fi = make_plain_file_info(Uuid::nil(), data_dir); + fi.version_id = None; + fi + } + fn make_file_info(version_id: Uuid, data_dir: Uuid) -> FileInfo { let restore_header = format!( "ongoing-request=\"false\", expiry-date=\"{}\"", @@ -113,17 +165,23 @@ mod tests { .format(&Rfc3339) .expect("format restore expiry"), ); + make_file_info_with_metadata( + version_id, + data_dir, + HashMap::from([ + ("etag".to_string(), format!("etag-{version_id}")), + (X_AMZ_RESTORE.as_str().to_string(), restore_header), + ]), + ) + } + + fn make_file_info_with_metadata(version_id: Uuid, data_dir: Uuid, metadata: HashMap) -> FileInfo { FileInfo { version_id: Some(version_id), data_dir: Some(data_dir), size: 64 * 1024, mod_time: Some(OffsetDateTime::now_utc()), - metadata: [ - ("etag".to_string(), format!("etag-{version_id}")), - (X_AMZ_RESTORE.as_str().to_string(), restore_header), - ] - .into_iter() - .collect(), + metadata, erasure: ErasureInfo { algorithm: ErasureAlgo::ReedSolomon.to_string(), data_blocks: 4, @@ -164,4 +222,80 @@ mod tests { let got = meta.find_unshared_data_dir_for_version(Some(target_version)); assert_eq!(got, None); } + + #[test] + fn find_unshared_data_dir_for_version_uses_data_dir_when_header_flag_is_false() { + let target_version = Uuid::new_v4(); + let target_data_dir = Uuid::new_v4(); + let mut meta = FileMeta::new(); + meta.add_version(make_plain_file_info(target_version, target_data_dir)) + .expect("seed target version"); + meta.add_version(make_plain_file_info(Uuid::new_v4(), Uuid::new_v4())) + .expect("seed non-shared version"); + + let target = meta + .versions + .iter() + .find(|version| version.header.version_id == Some(target_version)) + .expect("target version header"); + assert!(!target.header.uses_data_dir()); + + let got = meta.find_unshared_data_dir_for_version(Some(target_version)); + assert_eq!(got, Some(target_data_dir)); + } + + #[test] + fn find_unshared_data_dir_for_null_version_uses_data_dir_when_header_flag_is_false() { + let target_data_dir = Uuid::new_v4(); + let mut meta = FileMeta::new(); + meta.add_version(make_plain_null_version_file_info(target_data_dir)) + .expect("seed null version"); + + let target = meta + .versions + .iter() + .find(|version| version.header.version_id == Some(Uuid::nil())) + .expect("null version header"); + assert!(!target.header.uses_data_dir()); + + let got = meta.find_unshared_data_dir_for_version(None); + assert_eq!(got, Some(target_data_dir)); + } + + #[test] + fn shared_data_dir_count_uses_data_dir_when_header_flag_is_false() { + let target_version = Uuid::new_v4(); + let other_version = Uuid::new_v4(); + let shared_data_dir = Uuid::new_v4(); + let mut meta = FileMeta::new(); + meta.add_version(make_plain_file_info(target_version, shared_data_dir)) + .expect("seed target version"); + meta.add_version(make_plain_file_info(other_version, shared_data_dir)) + .expect("seed shared version"); + + assert!(meta.versions.iter().all(|version| !version.header.uses_data_dir())); + assert_eq!(meta.shared_data_dir_count(Some(target_version), Some(shared_data_dir)), 1); + + let old_dir = meta + .delete_version(&FileInfo { + version_id: Some(target_version), + ..Default::default() + }) + .expect("delete target version"); + assert_eq!(old_dir, None); + } + + #[test] + fn find_unshared_data_dir_for_version_skips_remote_transitioned_data_dir() { + let target_version = Uuid::new_v4(); + let target_data_dir = Uuid::new_v4(); + let mut fi = make_plain_file_info(target_version, target_data_dir); + fi.transition_status = TRANSITION_COMPLETE.to_string(); + + let mut meta = FileMeta::new(); + meta.add_version(fi).expect("seed remote transitioned version"); + + let got = meta.find_unshared_data_dir_for_version(Some(target_version)); + assert_eq!(got, None); + } }