fix(filemeta): detect physical data dirs for overwrite cleanup (#3510)

This commit is contained in:
GatewayJ
2026-06-18 11:58:12 +08:00
committed by GitHub
parent 24443c829d
commit a30cafa73f
3 changed files with 276 additions and 27 deletions
+3
View File
@@ -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;
@@ -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<dyn std::error::Error + Send + Sync>> {
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<Vec<Uuid>, Box<dyn std::error::Error + Send + Sync>> {
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)
}
+161 -27
View File
@@ -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<Option<Uuid>> {
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<Uuid>) -> Option<Uuid> {
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<Uuid>, data_dir: &Option<Uuid>) -> 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<Vec<Option<Uuid>>> {
@@ -81,6 +115,9 @@ impl FileMeta {
/// Count shared data directories
pub fn shared_data_dir_count(&self, version_id: Option<Uuid>, data_dir: Option<Uuid>) -> 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<String, String>) -> 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);
}
}