fix: filemeta version handling and delete operations (#879)

* fix filemeta version

* fix clippy

* fix delete version

* fix clippy/test
This commit is contained in:
weisd
2025-11-18 09:24:22 +08:00
committed by GitHub
parent 601f3456bc
commit 85bc0ce2d5
16 changed files with 93 additions and 96 deletions
+2 -2
View File
@@ -363,7 +363,7 @@ impl ErasureSetHealer {
let _permit = semaphore
.acquire()
.await
.map_err(|e| Error::other(format!("Failed to acquire semaphore for bucket heal: {}", e)))?;
.map_err(|e| Error::other(format!("Failed to acquire semaphore for bucket heal: {e}")))?;
if cancel_token.is_cancelled() {
return Err(Error::TaskCancelled);
@@ -461,7 +461,7 @@ impl ErasureSetHealer {
let _permit = semaphore
.acquire()
.await
.map_err(|e| Error::other(format!("Failed to acquire semaphore for object heal: {}", e)))?;
.map_err(|e| Error::other(format!("Failed to acquire semaphore for object heal: {e}")))?;
match storage.heal_object(&bucket, &object, None, &heal_opts).await {
Ok((_result, None)) => {
+1 -1
View File
@@ -30,7 +30,7 @@ const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json";
/// Helper function to convert Path to &str, returning an error if conversion fails
fn path_to_str(path: &Path) -> Result<&str> {
path.to_str()
.ok_or_else(|| Error::other(format!("Invalid UTF-8 path: {:?}", path)))
.ok_or_else(|| Error::other(format!("Invalid UTF-8 path: {path:?}")))
}
/// resume state
+1 -2
View File
@@ -180,8 +180,7 @@ impl HealStorageAPI for ECStoreHealStorage {
MAX_READ_BYTES, bucket, object
);
return Err(Error::other(format!(
"Object too large: {} bytes (max: {} bytes) for {}/{}",
n_read, MAX_READ_BYTES, bucket, object
"Object too large: {n_read} bytes (max: {MAX_READ_BYTES} bytes) for {bucket}/{object}"
)));
}
}
+6 -6
View File
@@ -144,7 +144,7 @@ async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>) {
let mut wtxn = lmdb_env.write_txn().unwrap();
let db = match lmdb_env
.database_options()
.name(&format!("bucket_{}", bucket_name))
.name(&format!("bucket_{bucket_name}"))
.types::<I64<BigEndian>, LifecycleContentCodec>()
.flags(DatabaseFlags::DUP_SORT)
//.dup_sort_comparator::<>()
@@ -152,7 +152,7 @@ async fn setup_test_env() -> (Vec<PathBuf>, Arc<ECStore>) {
{
Ok(db) => db,
Err(err) => {
panic!("lmdb error: {}", err);
panic!("lmdb error: {err}");
}
};
let _ = wtxn.commit();
@@ -199,7 +199,7 @@ async fn upload_test_object(ecstore: &Arc<ECStore>, bucket: &str, object: &str,
.await
.expect("Failed to upload test object");
println!("object_info1: {:?}", object_info);
println!("object_info1: {object_info:?}");
info!("Uploaded test object: {}/{} ({} bytes)", bucket, object, object_info.size);
}
@@ -456,7 +456,7 @@ mod serial_tests {
}
let object_info = convert_record_to_object_info(record);
println!("object_info2: {:?}", object_info);
println!("object_info2: {object_info:?}");
let mod_time = object_info.mod_time.unwrap_or(OffsetDateTime::now_utc());
let expiry_time = rustfs_ecstore::bucket::lifecycle::lifecycle::expected_expiry_time(mod_time, 1);
@@ -494,9 +494,9 @@ mod serial_tests {
type_,
object_name,
} = &elm.1;
println!("cache row:{} {} {} {:?} {}", ver_no, ver_id, mod_time, type_, object_name);
println!("cache row:{ver_no} {ver_id} {mod_time} {type_:?} {object_name}");
}
println!("row:{:?}", row);
println!("row:{row:?}");
}
//drop(iter);
wtxn.commit().unwrap();
@@ -277,11 +277,11 @@ async fn create_test_tier(server: u32) {
};
let mut tier_config_mgr = GLOBAL_TierConfigMgr.write().await;
if let Err(err) = tier_config_mgr.add(args, false).await {
println!("tier_config_mgr add failed, e: {:?}", err);
println!("tier_config_mgr add failed, e: {err:?}");
panic!("tier add failed. {err}");
}
if let Err(e) = tier_config_mgr.save().await {
println!("tier_config_mgr save failed, e: {:?}", e);
println!("tier_config_mgr save failed, e: {e:?}");
panic!("tier save failed");
}
println!("Created test tier: COLDTIER44");
@@ -299,7 +299,7 @@ async fn object_exists(ecstore: &Arc<ECStore>, bucket: &str, object: &str) -> bo
#[allow(dead_code)]
async fn object_is_delete_marker(ecstore: &Arc<ECStore>, bucket: &str, object: &str) -> bool {
if let Ok(oi) = (**ecstore).get_object_info(bucket, object, &ObjectOptions::default()).await {
println!("oi: {:?}", oi);
println!("oi: {oi:?}");
oi.delete_marker
} else {
println!("object_is_delete_marker is error");
@@ -311,7 +311,7 @@ async fn object_is_delete_marker(ecstore: &Arc<ECStore>, bucket: &str, object: &
#[allow(dead_code)]
async fn object_is_transitioned(ecstore: &Arc<ECStore>, bucket: &str, object: &str) -> bool {
if let Ok(oi) = (**ecstore).get_object_info(bucket, object, &ObjectOptions::default()).await {
println!("oi: {:?}", oi);
println!("oi: {oi:?}");
!oi.transitioned_object.status.is_empty()
} else {
println!("object_is_transitioned is error");
@@ -1106,18 +1106,15 @@ impl TargetClient {
SdkError::ServiceError(oe) => match oe.into_err() {
HeadBucketError::NotFound(_) => Ok(false),
other => Err(S3ClientError::new(format!(
"failed to check bucket exists for bucket:{} please check the bucket name and credentials, error:{:?}",
bucket, other
"failed to check bucket exists for bucket:{bucket} please check the bucket name and credentials, error:{other:?}"
))),
},
SdkError::DispatchFailure(e) => Err(S3ClientError::new(format!(
"failed to dispatch bucket exists for bucket:{} error:{:?}",
bucket, e
"failed to dispatch bucket exists for bucket:{bucket} error:{e:?}"
))),
_ => Err(S3ClientError::new(format!(
"failed to check bucket exists for bucket:{} error:{:?}",
bucket, e
"failed to check bucket exists for bucket:{bucket} error:{e:?}"
))),
},
}
+2 -3
View File
@@ -806,7 +806,7 @@ impl LocalDisk {
Ok((bytes, modtime))
}
async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &Vec<FileInfo>) -> Result<()> {
async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &[FileInfo]) -> Result<()> {
let volume_dir = self.get_bucket_path(volume)?;
let xlpath = self.get_object_path(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str())?;
@@ -820,7 +820,7 @@ impl LocalDisk {
fm.unmarshal_msg(&data)?;
for fi in fis {
for fi in fis.iter() {
let data_dir = match fm.delete_version(fi) {
Ok(res) => res,
Err(err) => {
@@ -2300,7 +2300,6 @@ impl DiskAPI for LocalDisk {
let buf = match self.read_all_data(volume, &volume_dir, &xl_path).await {
Ok(res) => res,
Err(err) => {
//
if err != DiskError::FileNotFound {
return Err(err);
}
+1 -1
View File
@@ -1085,7 +1085,7 @@ mod tests {
drop(listener);
let url = url::Url::parse(&format!("http://{}:{}/data/rustfs0", ip, port)).unwrap();
let url = url::Url::parse(&format!("http://{ip}:{port}/data/rustfs0")).unwrap();
let endpoint = Endpoint {
url,
is_local: false,
+1 -2
View File
@@ -4247,7 +4247,6 @@ impl StorageAPI for SetDisks {
for (_, mut fi_vers) in vers_map {
fi_vers.versions.sort_by(|a, b| a.deleted.cmp(&b.deleted));
fi_vers.versions.reverse();
if let Some(index) = fi_vers.versions.iter().position(|fi| fi.deleted) {
fi_vers.versions.truncate(index + 1);
@@ -4652,7 +4651,7 @@ impl StorageAPI for SetDisks {
let tgt_client = match tier_config_mgr.get_driver(&opts.transition.tier).await {
Ok(client) => client,
Err(err) => {
return Err(Error::other(format!("remote tier error: {}", err)));
return Err(Error::other(format!("remote tier error: {err}")));
}
};
+41 -29
View File
@@ -427,11 +427,19 @@ impl FileMeta {
return;
}
self.versions.reverse();
// for _v in self.versions.iter() {
// // warn!("sort {} {:?}", i, v);
// }
self.versions.sort_by(|a, b| {
if a.header.mod_time != b.header.mod_time {
b.header.mod_time.cmp(&a.header.mod_time)
} else if a.header.version_type != b.header.version_type {
b.header.version_type.cmp(&a.header.version_type)
} else if a.header.version_id != b.header.version_id {
b.header.version_id.cmp(&a.header.version_id)
} else if a.header.flags != b.header.flags {
b.header.flags.cmp(&a.header.flags)
} else {
b.cmp(a)
}
});
}
// Find version
@@ -489,25 +497,27 @@ impl FileMeta {
self.versions.sort_by(|a, b| {
if a.header.mod_time != b.header.mod_time {
a.header.mod_time.cmp(&b.header.mod_time)
b.header.mod_time.cmp(&a.header.mod_time)
} else if a.header.version_type != b.header.version_type {
a.header.version_type.cmp(&b.header.version_type)
b.header.version_type.cmp(&a.header.version_type)
} else if a.header.version_id != b.header.version_id {
a.header.version_id.cmp(&b.header.version_id)
b.header.version_id.cmp(&a.header.version_id)
} else if a.header.flags != b.header.flags {
a.header.flags.cmp(&b.header.flags)
b.header.flags.cmp(&a.header.flags)
} else {
a.cmp(b)
b.cmp(a)
}
});
Ok(())
}
pub fn add_version(&mut self, fi: FileInfo) -> Result<()> {
let vid = fi.version_id;
pub fn add_version(&mut self, mut fi: FileInfo) -> Result<()> {
if fi.version_id.is_none() {
fi.version_id = Some(Uuid::nil());
}
if let Some(ref data) = fi.data {
let key = vid.unwrap_or_default().to_string();
let key = fi.version_id.unwrap_or_default().to_string();
self.data.replace(&key, data.to_vec())?;
}
@@ -551,7 +561,6 @@ impl FileMeta {
}
}
}
Err(Error::other("add_version failed"))
// if !ver.valid() {
@@ -583,12 +592,19 @@ impl FileMeta {
}
// delete_version deletes version, returns data_dir
#[tracing::instrument(skip(self))]
pub fn delete_version(&mut self, fi: &FileInfo) -> Result<Option<Uuid>> {
let vid = if fi.version_id.is_none() {
Some(Uuid::nil())
} else {
Some(fi.version_id.unwrap())
};
let mut ventry = FileMetaVersion::default();
if fi.deleted {
ventry.version_type = VersionType::Delete;
ventry.delete_marker = Some(MetaDeleteMarker {
version_id: fi.version_id,
version_id: vid,
mod_time: fi.mod_time,
..Default::default()
});
@@ -689,8 +705,10 @@ impl FileMeta {
}
}
let mut found_index = None;
for (i, ver) in self.versions.iter().enumerate() {
if ver.header.version_id != fi.version_id {
if ver.header.version_id != vid {
continue;
}
@@ -701,7 +719,7 @@ impl FileMeta {
let mut v = self.get_idx(i)?;
if v.delete_marker.is_none() {
v.delete_marker = Some(MetaDeleteMarker {
version_id: fi.version_id,
version_id: vid,
mod_time: fi.mod_time,
meta_sys: HashMap::new(),
});
@@ -767,7 +785,7 @@ impl FileMeta {
self.versions.remove(i);
if (fi.mark_deleted && fi.version_purge_status() != VersionPurgeStatusType::Complete)
|| (fi.deleted && fi.version_id.is_none())
|| (fi.deleted && vid == Some(Uuid::nil()))
{
self.add_version_filemata(ventry)?;
}
@@ -803,18 +821,11 @@ impl FileMeta {
return Ok(old_dir);
}
found_index = Some(i);
}
}
}
let mut found_index = None;
for (i, version) in self.versions.iter().enumerate() {
if version.header.version_type == VersionType::Object && version.header.version_id == fi.version_id {
found_index = Some(i);
break;
}
}
let Some(i) = found_index else {
if fi.deleted {
self.add_version_filemata(ventry)?;
@@ -1521,7 +1532,8 @@ impl FileMetaVersionHeader {
cur.read_exact(&mut buf)?;
self.version_id = {
let id = Uuid::from_bytes(buf);
if id.is_nil() { None } else { Some(id) }
// if id.is_nil() { None } else { Some(id) }
Some(id)
};
// mod_time
@@ -1695,7 +1707,7 @@ impl MetaObject {
}
pub fn into_fileinfo(&self, volume: &str, path: &str, all_parts: bool) -> FileInfo {
let version_id = self.version_id.filter(|&vid| !vid.is_nil());
// let version_id = self.version_id.filter(|&vid| !vid.is_nil());
let parts = if all_parts {
let mut parts = vec![ObjectPartInfo::default(); self.part_numbers.len()];
@@ -1799,7 +1811,7 @@ impl MetaObject {
.unwrap_or_default();
FileInfo {
version_id,
version_id: self.version_id,
erasure,
data_dir: self.data_dir,
mod_time: self.mod_time,
+3 -5
View File
@@ -272,8 +272,7 @@ fn init_file_logging(config: &OtelConfig, logger_level: &str, is_production: boo
if (current & !desired) != 0 {
if let Err(e) = fs::set_permissions(log_directory, Permissions::from_mode(desired)) {
return Err(TelemetryError::SetPermissions(format!(
"dir='{}', want={:#o}, have={:#o}, err={}",
log_directory, desired, current, e
"dir='{log_directory}', want={desired:#o}, have={current:#o}, err={e}"
)));
}
// Second verification
@@ -281,15 +280,14 @@ fn init_file_logging(config: &OtelConfig, logger_level: &str, is_production: boo
let after = meta2.permissions().mode() & 0o777;
if after != desired {
return Err(TelemetryError::SetPermissions(format!(
"dir='{}', want={:#o}, after={:#o}",
log_directory, desired, after
"dir='{log_directory}', want={desired:#o}, after={after:#o}"
)));
}
}
}
}
Err(e) => {
return Err(TelemetryError::Io(format!("stat '{}' failed: {}", log_directory, e)));
return Err(TelemetryError::Io(format!("stat '{log_directory}' failed: {e}")));
}
}
}
+9 -9
View File
@@ -543,7 +543,7 @@ mod tests {
*req.uri_mut() = Uri::from_parts(parts).unwrap();
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
println!("canonical_request: \n{}\n", canonical_request);
println!("canonical_request: \n{canonical_request}\n");
assert_eq!(
canonical_request,
concat!(
@@ -562,7 +562,7 @@ mod tests {
);
let string_to_sign = get_string_to_sign_v4(t, region, &canonical_request, service);
println!("string_to_sign: \n{}\n", string_to_sign);
println!("string_to_sign: \n{string_to_sign}\n");
assert_eq!(
string_to_sign,
concat!(
@@ -576,7 +576,7 @@ mod tests {
let signing_key = get_signing_key(secret_access_key, region, t, service);
let signature = get_signature(signing_key, &string_to_sign);
println!("signature: \n{}\n", signature);
println!("signature: \n{signature}\n");
assert_eq!(signature, "73fad2dfea0727e10a7179bf49150360a56f2e6b519c53999fd6e011152187d0");
}
@@ -608,7 +608,7 @@ mod tests {
println!("{:?}", req.uri().query());
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
println!("canonical_request: \n{}\n", canonical_request);
println!("canonical_request: \n{canonical_request}\n");
assert_eq!(
canonical_request,
concat!(
@@ -627,7 +627,7 @@ mod tests {
);
let string_to_sign = get_string_to_sign_v4(t, region, &canonical_request, service);
println!("string_to_sign: \n{}\n", string_to_sign);
println!("string_to_sign: \n{string_to_sign}\n");
assert_eq!(
string_to_sign,
concat!(
@@ -641,7 +641,7 @@ mod tests {
let signing_key = get_signing_key(secret_access_key, region, t, service);
let signature = get_signature(signing_key, &string_to_sign);
println!("signature: \n{}\n", signature);
println!("signature: \n{signature}\n");
assert_eq!(signature, "dfbed913d1982428f6224ee506431fc133dbcad184194c0cbf01bc517435788a");
}
@@ -673,7 +673,7 @@ mod tests {
println!("{:?}", req.uri().query());
let canonical_request = get_canonical_request(&req, &v4_ignored_headers, &get_hashed_payload(&req));
println!("canonical_request: \n{}\n", canonical_request);
println!("canonical_request: \n{canonical_request}\n");
assert_eq!(
canonical_request,
concat!(
@@ -692,7 +692,7 @@ mod tests {
);
let string_to_sign = get_string_to_sign_v4(t, region, &canonical_request, service);
println!("string_to_sign: \n{}\n", string_to_sign);
println!("string_to_sign: \n{string_to_sign}\n");
assert_eq!(
string_to_sign,
concat!(
@@ -706,7 +706,7 @@ mod tests {
let signing_key = get_signing_key(secret_access_key, region, t, service);
let signature = get_signature(signing_key, &string_to_sign);
println!("signature: \n{}\n", signature);
println!("signature: \n{signature}\n");
assert_eq!(signature, "c7c7c6e12e5709c0c2ffc4707600a86c3cd261dd1de7409126a17f5b08c58dfa");
}
+1 -1
View File
@@ -114,7 +114,7 @@ async fn async_main() -> Result<()> {
let guard = match init_obs(Some(opt.clone().obs_endpoint)).await {
Ok(g) => g,
Err(e) => {
println!("Failed to initialize observability: {}", e);
println!("Failed to initialize observability: {e}");
return Err(Error::other(e));
}
};
+2 -2
View File
@@ -1077,7 +1077,7 @@ impl S3 for FS {
warn!("unable to restore transitioned bucket/object {}/{}: {}", bucket, object, err.to_string());
return Err(S3Error::with_message(
S3ErrorCode::Custom("ErrRestoreTransitionedObject".into()),
format!("unable to restore transitioned bucket/object {}/{}: {}", bucket, object, err),
format!("unable to restore transitioned bucket/object {bucket}/{object}: {err}"),
));
}
@@ -1329,7 +1329,7 @@ impl S3 for FS {
}
if is_dir_object(&object.object_name) && object.version_id.is_none() {
object.version_id = Some(Uuid::max());
object.version_id = Some(Uuid::nil());
}
if replicate_deletes {
+1 -1
View File
@@ -157,7 +157,7 @@ impl OperationHelper {
.clone()
.status(status)
.status_code(status_code)
.time_to_response(format!("{:.2?}", ttr))
.time_to_response(format!("{ttr:.2?}"))
.time_to_response_in_ns(ttr.as_nanos().to_string())
.build();
+15 -22
View File
@@ -65,15 +65,12 @@ pub async fn del_opts(
let vid = vid.map(|v| v.as_str().trim().to_owned());
if let Some(ref id) = vid {
if let Err(err) = Uuid::parse_str(id.as_str()) {
if *id != Uuid::nil().to_string()
&& let Err(err) = Uuid::parse_str(id.as_str())
{
error!("del_opts: invalid version id: {} error: {}", id, err);
return Err(StorageError::InvalidVersionID(bucket.to_owned(), object.to_owned(), id.clone()));
}
if !versioned {
error!("del_opts: object not versioned: {}", object);
return Err(StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), id.clone()));
}
}
let mut opts = put_opts_from_headers(headers, metadata.clone()).map_err(|err| {
@@ -83,7 +80,7 @@ pub async fn del_opts(
opts.version_id = {
if is_dir_object(object) && vid.is_none() {
Some(Uuid::max().to_string())
Some(Uuid::nil().to_string())
} else {
vid
}
@@ -113,13 +110,11 @@ pub async fn get_opts(
let vid = vid.map(|v| v.as_str().trim().to_owned());
if let Some(ref id) = vid {
if let Err(_err) = Uuid::parse_str(id.as_str()) {
if *id != Uuid::nil().to_string()
&& let Err(_err) = Uuid::parse_str(id.as_str())
{
return Err(StorageError::InvalidVersionID(bucket.to_owned(), object.to_owned(), id.clone()));
}
if !versioned {
return Err(StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), id.clone()));
}
}
let mut opts = get_default_opts(headers, HashMap::new(), false)
@@ -127,7 +122,7 @@ pub async fn get_opts(
opts.version_id = {
if is_dir_object(object) && vid.is_none() {
Some(Uuid::max().to_string())
Some(Uuid::nil().to_string())
} else {
vid
}
@@ -189,13 +184,11 @@ pub async fn put_opts(
let vid = vid.map(|v| v.as_str().trim().to_owned());
if let Some(ref id) = vid {
if let Err(_err) = Uuid::parse_str(id.as_str()) {
if *id != Uuid::nil().to_string()
&& let Err(_err) = Uuid::parse_str(id.as_str())
{
return Err(StorageError::InvalidVersionID(bucket.to_owned(), object.to_owned(), id.clone()));
}
if !versioned {
return Err(StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), id.clone()));
}
}
let mut opts = put_opts_from_headers(headers, metadata)
@@ -203,7 +196,7 @@ pub async fn put_opts(
opts.version_id = {
if is_dir_object(object) && vid.is_none() {
Some(Uuid::max().to_string())
Some(Uuid::nil().to_string())
} else {
vid
}
@@ -603,7 +596,7 @@ mod tests {
assert!(result.is_ok());
let opts = result.unwrap();
assert_eq!(opts.version_id, Some(Uuid::max().to_string()));
assert_eq!(opts.version_id, Some(Uuid::nil().to_string()));
}
#[tokio::test]
@@ -676,7 +669,7 @@ mod tests {
assert!(result.is_ok());
let opts = result.unwrap();
assert_eq!(opts.version_id, Some(Uuid::max().to_string()));
assert_eq!(opts.version_id, Some(Uuid::nil().to_string()));
}
#[tokio::test]
@@ -720,7 +713,7 @@ mod tests {
assert!(result.is_ok());
let opts = result.unwrap();
assert_eq!(opts.version_id, Some(Uuid::max().to_string()));
assert_eq!(opts.version_id, Some(Uuid::nil().to_string()));
}
#[tokio::test]