From bd36cf358880d69566eff24bdb438b0b007e5300 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Tue, 31 Mar 2026 22:08:45 +0800 Subject: [PATCH] test(filemeta): cover legacy delete marker decoding (#2333) Signed-off-by: dependabot[bot] Co-authored-by: houseme Co-authored-by: heihutu Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> --- crates/filemeta/src/filemeta/version.rs | 44 +++++++++++++++++++++++ crates/io-metrics/src/capacity_metrics.rs | 26 +++++++++----- crates/s3-common/src/event_name.rs | 11 +++--- rustfs/src/app/object_usecase.rs | 29 ++++++++++++--- rustfs/src/capacity/capacity_manager.rs | 12 +++++-- 5 files changed, 102 insertions(+), 20 deletions(-) diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index 359cbff54..06b609545 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -2966,6 +2966,38 @@ mod tests { assert!(fi.uses_legacy_checksum); } + #[test] + fn legacy_meta_v2_delete_marker_decodes_into_delete_fileinfo_via_struct() { + let version_id = sample_version_id(); + let mod_time = sample_mod_time(); + let version = LegacyMetaV2Version { + version_type: LegacyMetaV2VersionType::DeleteMarker, + object: None, + delete_marker: Some(LegacyMetaV2DeleteMarker { + version_id: version_id.as_bytes().to_vec(), + mod_time: Some(mod_time), + meta_sys: HashMap::from([("x-minio-internal".to_string(), b"present".to_vec())]), + }), + write_version: 7, + }; + + let decoded = FileMetaVersion::try_from(version).unwrap(); + + assert_eq!(decoded.version_type, VersionType::Delete); + assert!(decoded.uses_legacy_checksum); + assert!(decoded.object.is_none()); + + let delete_marker = decoded.delete_marker.as_ref().expect("delete marker should be decoded"); + assert_eq!(delete_marker.version_id, Some(version_id)); + assert_eq!(delete_marker.mod_time, Some(mod_time)); + + let fi = decoded.into_fileinfo("bucket", "deleted.txt", true); + assert!(fi.deleted); + assert_eq!(fi.version_id, Some(version_id)); + assert_eq!(fi.mod_time, Some(mod_time)); + assert_eq!(fi.metadata.get("x-minio-internal").map(String::as_str), Some("present")); + } + #[test] fn legacy_meta_v2_delete_marker_rejects_invalid_uuid_bytes() { let payload = LegacyDeleteVersionFixture { @@ -2983,4 +3015,16 @@ mod tests { let err = FileMetaVersion::try_from(encoded.as_slice()).expect_err("invalid legacy delete marker UUID must fail"); assert!(err.to_string().contains("legacy version_id must be 16 bytes")); } + + #[test] + fn legacy_meta_v2_delete_marker_rejects_invalid_uuid_bytes_via_struct() { + let err = MetaDeleteMarker::try_from(LegacyMetaV2DeleteMarker { + version_id: vec![1, 2, 3], + mod_time: Some(sample_mod_time()), + meta_sys: HashMap::new(), + }) + .expect_err("invalid legacy delete-marker version ids should be rejected"); + + assert!(err.to_string().contains("legacy version_id must be 16 bytes")); + } } diff --git a/crates/io-metrics/src/capacity_metrics.rs b/crates/io-metrics/src/capacity_metrics.rs index 070d67cc9..a032727ee 100644 --- a/crates/io-metrics/src/capacity_metrics.rs +++ b/crates/io-metrics/src/capacity_metrics.rs @@ -37,18 +37,22 @@ pub fn record_capacity_current_bytes(used_bytes: u64) { /// Record capacity update completion. #[inline(always)] -pub fn record_capacity_update_completed(source: &str, duration: Duration, used_bytes: u64, is_estimated: bool) { - counter!("rustfs.capacity.update.total", "source" => source.to_string()).increment(1); - histogram!("rustfs.capacity.update.duration.seconds", "source" => source.to_string()).record(duration.as_secs_f64()); - histogram!("rustfs.capacity.update.bytes", "source" => source.to_string()).record(used_bytes as f64); - counter!("rustfs.capacity.update.estimated.total", "source" => source.to_string(), "estimated" => is_estimated.to_string()) - .increment(1); +pub fn record_capacity_update_completed(source: &'static str, duration: Duration, used_bytes: u64, is_estimated: bool) { + counter!("rustfs.capacity.update.total", "source" => source).increment(1); + histogram!("rustfs.capacity.update.duration.seconds", "source" => source).record(duration.as_secs_f64()); + histogram!("rustfs.capacity.update.bytes", "source" => source).record(used_bytes as f64); + counter!( + "rustfs.capacity.update.estimated.total", + "source" => source, + "estimated" => if is_estimated { "true" } else { "false" } + ) + .increment(1); } /// Record failed capacity update. #[inline(always)] -pub fn record_capacity_update_failed(source: &str) { - counter!("rustfs.capacity.update.failures", "source" => source.to_string()).increment(1); +pub fn record_capacity_update_failed(source: &'static str) { + counter!("rustfs.capacity.update.failures", "source" => source).increment(1); } /// Record capacity write activity. @@ -88,5 +92,9 @@ pub fn record_capacity_dynamic_timeout(timeout: Duration) { #[inline(always)] pub fn record_capacity_scan_sampling(sampled_count: usize, estimated: bool) { histogram!("rustfs.capacity.scan.sampled.count").record(sampled_count as f64); - counter!("rustfs.capacity.scan.estimated.total", "estimated" => estimated.to_string()).increment(1); + counter!( + "rustfs.capacity.scan.estimated.total", + "estimated" => if estimated { "true" } else { "false" } + ) + .increment(1); } diff --git a/crates/s3-common/src/event_name.rs b/crates/s3-common/src/event_name.rs index 6b22a656e..7776dc751 100644 --- a/crates/s3-common/src/event_name.rs +++ b/crates/s3-common/src/event_name.rs @@ -164,7 +164,7 @@ impl EventName { "s3:ObjectRemoved:DeleteMarkerCreated" => Ok(EventName::ObjectRemovedDeleteMarkerCreated), "s3:ObjectRemoved:NoOP" => Ok(EventName::ObjectRemovedNoOP), "s3:ObjectRemoved:DeleteAllVersions" => Ok(EventName::ObjectRemovedDeleteAllVersions), - "s3:LifecycleDelMarkerExpiration:Delete" => Ok(EventName::LifecycleExpirationDeleteMarkerCreated), + "s3:LifecycleDelMarkerExpiration:Delete" => Ok(EventName::LifecycleDelMarkerExpirationDelete), "s3:LifecycleExpiration:*" => Ok(EventName::LifecycleExpirationAll), "s3:LifecycleExpiration:Delete" => Ok(EventName::LifecycleExpirationDelete), "s3:LifecycleExpiration:DeleteMarkerCreated" => Ok(EventName::LifecycleExpirationDeleteMarkerCreated), @@ -178,7 +178,7 @@ impl EventName { "s3:ObjectRestore:Post" => Ok(EventName::ObjectRestorePost), "s3:ObjectRestore:Completed" => Ok(EventName::ObjectRestoreCompleted), "s3:ObjectTransition:Failed" => Ok(EventName::ObjectTransitionFailed), - "s3:ObjectTransition:Complete" => Ok(EventName::LifecycleTransition), + "s3:ObjectTransition:Complete" => Ok(EventName::ObjectTransitionComplete), "s3:ObjectTransition:*" => Ok(EventName::ObjectTransitionAll), "s3:LifecycleTransition" => Ok(EventName::LifecycleTransition), "s3:IntelligentTiering" => Ok(EventName::IntelligentTiering), @@ -650,10 +650,13 @@ mod tests { EventName::parse("s3:ObjectCreated:DeleteTagging").unwrap(), EventName::ObjectTaggingDelete ); - assert_eq!(EventName::parse("s3:ObjectTransition:Complete").unwrap(), EventName::LifecycleTransition); + assert_eq!( + EventName::parse("s3:ObjectTransition:Complete").unwrap(), + EventName::ObjectTransitionComplete + ); assert_eq!( EventName::parse("s3:LifecycleDelMarkerExpiration:Delete").unwrap(), - EventName::LifecycleExpirationDeleteMarkerCreated + EventName::LifecycleDelMarkerExpirationDelete ); } diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 8ad0bf011..b8c70de93 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -1631,13 +1631,13 @@ impl DefaultObjectUsecase { #[instrument(level = "debug", skip(self, _fs, req))] pub async fn execute_put_object(&self, _fs: &FS, req: S3Request) -> S3Result> { let start_time = std::time::Instant::now(); + let mut req = req; if let Some(context) = &self.context { let _ = context.object_store(); } let (event_name, quota_operation, request_method_name) = Self::put_object_execution_context(&req); - let mut helper = OperationHelper::new(&req, event_name, S3Operation::PutObject); if req.extensions.get::().is_some() && is_post_object_sse_kms_requested(&req.input, &req.headers) { return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for POST object uploads")); @@ -1651,7 +1651,7 @@ impl DefaultObjectUsecase { return self.execute_put_object_extract(req).await; } - let input = req.input; + let input = std::mem::take(&mut req.input); let PutObjectInput { body, @@ -1870,6 +1870,8 @@ impl DefaultObjectUsecase { opts.want_checksum = reader.checksum(); } + let mut helper = OperationHelper::new(&req, event_name, S3Operation::PutObject); + // Apply encryption using unified SSE API. let encryption_request = EncryptionRequest { bucket: &bucket, @@ -1885,7 +1887,16 @@ impl DefaultObjectUsecase { part_nonce: None, }; - if let Some(material) = sse_encryption(encryption_request).await? { + let encryption_material = match sse_encryption(encryption_request).await { + Ok(material) => material, + Err(err) => { + let result = Err(err.into()); + let _ = helper.complete(&result); + return result; + } + }; + + if let Some(material) = encryption_material { effective_sse = Some(material.server_side_encryption.clone()); effective_kms_key_id = material.kms_key_id.clone(); @@ -1917,10 +1928,18 @@ impl DefaultObjectUsecase { ); } - let obj_info = store + let obj_info = match store .put_object(&bucket, &key, &mut reader, &opts) .await - .map_err(ApiError::from)?; + .map_err(ApiError::from) + { + Ok(obj_info) => obj_info, + Err(err) => { + let result: S3Result> = Err(err.into()); + let _ = helper.complete(&result); + return result; + } + }; maybe_enqueue_transition_immediate(&obj_info, LcEventSrc::S3PutObject).await; diff --git a/rustfs/src/capacity/capacity_manager.rs b/rustfs/src/capacity/capacity_manager.rs index 8f1096215..d8586770f 100644 --- a/rustfs/src/capacity/capacity_manager.rs +++ b/rustfs/src/capacity/capacity_manager.rs @@ -15,6 +15,7 @@ //! Hybrid Capacity Manager for efficient capacity statistics use crate::app::admin_usecase::calculate_data_dir_used_capacity; +use futures::FutureExt; use rustfs_config::{ DEFAULT_CAPACITY_ENABLE_DYNAMIC_TIMEOUT, DEFAULT_CAPACITY_FOLLOW_SYMLINKS, DEFAULT_CAPACITY_MAX_SYMLINK_DEPTH, DEFAULT_CAPACITY_MAX_TIMEOUT_SECS, DEFAULT_CAPACITY_MIN_TIMEOUT_SECS, DEFAULT_CAPACITY_STALL_TIMEOUT_SECS, @@ -28,6 +29,7 @@ use rustfs_config::{ use rustfs_io_metrics::{record_capacity_current_bytes, record_capacity_update_completed, record_capacity_write_operation}; use rustfs_utils::{get_env_bool, get_env_u64, get_env_usize}; use std::future::Future; +use std::panic::AssertUnwindSafe; use std::sync::Arc; use std::time::{Duration, Instant}; use tokio::sync::{Mutex, RwLock, watch}; @@ -601,7 +603,10 @@ impl HybridCapacityManager { .unwrap_or_else(|| Err("capacity refresh completed without a result".to_string())); } - let result = refresh_fn().await; + let result = AssertUnwindSafe(refresh_fn()).catch_unwind().await.unwrap_or_else(|err| { + warn!(error = ?err, "capacity refresh function panicked"); + Err("capacity refresh panicked".to_string()) + }); if let Ok(update) = &result { self.update_capacity(update.clone(), source).await; } @@ -638,7 +643,10 @@ impl HybridCapacityManager { } tokio::spawn(async move { - let result = refresh_fn().await; + let result = AssertUnwindSafe(refresh_fn()).catch_unwind().await.unwrap_or_else(|err| { + warn!(error = ?err, "capacity refresh function panicked"); + Err("capacity refresh panicked".to_string()) + }); if let Ok(update) = &result { self.update_capacity(update.clone(), source).await; }