From 5ef8b1ce5c5ed930f9d6c8b289e01c14ee54f5a6 Mon Sep 17 00:00:00 2001 From: houseme Date: Fri, 28 Aug 2026 22:46:00 +0800 Subject: [PATCH] fix(tier): harden reference proof and audit output (#6807) Validate lifecycle tier references through the tier reference proof path, preserve S3 list CommonPrefix XML compatibility, and make GetObject audit completion use real S3 error status codes. Co-authored-by: heihutu --- crates/ecstore/src/bucket/lifecycle/mod.rs | 2 +- crates/ecstore/src/services/tier/tier.rs | 222 +++++++++++++++++++-- crates/s3-client/src/api_s3_datatypes.rs | 27 +++ rustfs/src/app/object/get.rs | 29 ++- rustfs/src/storage/helper.rs | 75 ++++++- 5 files changed, 326 insertions(+), 29 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/mod.rs b/crates/ecstore/src/bucket/lifecycle/mod.rs index d87e5436f..f8e286e8f 100644 --- a/crates/ecstore/src/bucket/lifecycle/mod.rs +++ b/crates/ecstore/src/bucket/lifecycle/mod.rs @@ -20,7 +20,7 @@ mod durable_namespace; pub mod evaluator; pub mod manual_transition_job; mod metadata_boundary; -pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs}; +pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs, get_lifecycle_config}; mod object_handlers_common; mod object_lock_boundary; pub use self::core as lifecycle; diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index bac56bd0b..1024f1791 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -65,6 +65,7 @@ use crate::storage_api_contracts::{ }; use crate::{ bucket::lifecycle::{ + get_lifecycle_config, tier_delete_journal::{TIER_DELETE_JOURNAL_PREFIX, decode_tier_delete_journal_entry}, transition_transaction::{TRANSITION_TRANSACTION_RECORD_PREFIX, decode_transition_transaction_record}, }, @@ -81,7 +82,7 @@ use rustfs_filemeta::FileInfo; use rustfs_rio::HashReader; use rustfs_s3_client::{admin_handler_utils::AdminError, provider_versions::ProviderVersionCapabilities}; use rustfs_utils::path::{SLASH_SEPARATOR, path_join}; -use s3s::S3ErrorCode; +use s3s::{S3ErrorCode, dto::BucketLifecycleConfiguration}; use super::{ tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_PERM_ERR}, @@ -694,6 +695,7 @@ fn tier_backend_identity_admin_error(err: io::Error) -> AdminError { admin_err } +#[async_trait::async_trait] trait TierReferenceProofStore: EcstoreObjectIO + BucketOperations @@ -705,23 +707,21 @@ trait TierReferenceProofStore: WalkOptions = TierReferenceProofWalkOptions, WalkCancellation = tokio_util::sync::CancellationToken, WalkResultSender = tokio::sync::mpsc::Sender>, - > + > + Send + + Sync { + async fn lifecycle_config_for_reference_proof(&self, bucket: &str) -> Result>; } -impl TierReferenceProofStore for T where - T: EcstoreObjectIO - + BucketOperations - + ListOperations< - Error = Error, - ListObjectsV2Info = StorageListObjectsV2Info, - ListObjectVersionsInfo = StorageListObjectVersionsInfo, - ObjectInfoOrErr = StorageObjectInfoOrErr, - WalkOptions = TierReferenceProofWalkOptions, - WalkCancellation = tokio_util::sync::CancellationToken, - WalkResultSender = tokio::sync::mpsc::Sender>, - > -{ +#[async_trait::async_trait] +impl TierReferenceProofStore for ECStore { + async fn lifecycle_config_for_reference_proof(&self, bucket: &str) -> Result> { + match get_lifecycle_config(bucket).await { + Ok((config, _updated_at)) => Ok(Some(config)), + Err(Error::ConfigNotFound) => Ok(None), + Err(err) => Err(err), + } + } } async fn ensure_no_authoritative_tier_object_references( @@ -754,6 +754,7 @@ where .await .map_err(tier_reference_proof_admin_error)?; for bucket in buckets { + ensure_no_authoritative_lifecycle_references(api.as_ref(), &bucket.name, targets).await?; let mut marker = None; let mut version_marker = None; loop { @@ -789,6 +790,66 @@ where ensure_no_authoritative_persisted_references(api, targets).await } +async fn ensure_no_authoritative_lifecycle_references( + api: &impl TierReferenceProofStore, + bucket: &str, + targets: &[TierMutationIntentTarget], +) -> std::result::Result<(), AdminError> { + match api.lifecycle_config_for_reference_proof(bucket).await { + Ok(Some(config)) => { + if let Some(reference) = lifecycle_config_target_reference(&config, targets) { + return Err(tier_reference_proof_lifecycle_in_use_error( + reference.tier_name, + bucket, + reference.rule_id, + )); + } + } + Ok(None) => {} + Err(err) => { + return Err(tier_reference_proof_admin_error(format!("bucket {bucket} lifecycle config: {err}"))); + } + } + Ok(()) +} + +struct TierLifecycleReference<'a> { + tier_name: &'a str, + rule_id: Option<&'a str>, +} + +fn lifecycle_config_target_reference<'a>( + config: &'a BucketLifecycleConfiguration, + targets: &[TierMutationIntentTarget], +) -> Option> { + for rule in &config.rules { + let rule_id = rule.id.as_deref(); + if let Some(transitions) = &rule.transitions { + for transition in transitions { + let Some(storage_class) = &transition.storage_class else { + continue; + }; + let tier_name = storage_class.as_str(); + if !tier_name.is_empty() && targets.iter().any(|target| target.tier_name == tier_name) { + return Some(TierLifecycleReference { tier_name, rule_id }); + } + } + } + if let Some(noncurrent_version_transitions) = &rule.noncurrent_version_transitions { + for transition in noncurrent_version_transitions { + let Some(storage_class) = &transition.storage_class else { + continue; + }; + let tier_name = storage_class.as_str(); + if !tier_name.is_empty() && targets.iter().any(|target| target.tier_name == tier_name) { + return Some(TierLifecycleReference { tier_name, rule_id }); + } + } + } + } + None +} + async fn ensure_no_authoritative_persisted_references( api: Arc, targets: &[TierMutationIntentTarget], @@ -920,6 +981,15 @@ fn tier_reference_proof_persisted_in_use_error(tier_name: &str, object: &str) -> err } +fn tier_reference_proof_lifecycle_in_use_error(tier_name: &str, bucket: &str, rule_id: Option<&str>) -> AdminError { + let mut err = ERR_TIER_BACKEND_IN_USE.clone(); + err.message = match rule_id { + Some(rule_id) => format!("Remote tier {tier_name} is still referenced by lifecycle rule {rule_id} in bucket {bucket}"), + None => format!("Remote tier {tier_name} is still referenced by a lifecycle rule in bucket {bucket}"), + }; + err +} + fn tier_reference_proof_admin_error(err: impl std::fmt::Display) -> AdminError { let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); admin_err.message = format!("Remote tier reference proof failed: {err}"); @@ -5308,6 +5378,10 @@ mod tests { endpoints::{Endpoints, PoolEndpoints, SetupType}, }; use crate::services::tier::tier_mutation_intent::TIER_MUTATION_INTENT_RECORD_PREFIX; + use s3s::dto::{ + BucketLifecycleConfiguration, ExpirationStatus, LifecycleRule, NoncurrentVersionTransition, Transition, + TransitionStorageClass, + }; struct SetupTypeGuard { previous: SetupType, @@ -6093,6 +6167,13 @@ mod tests { } } + #[async_trait::async_trait] + impl TierReferenceProofStore for LockingTierConfigStore { + async fn lifecycle_config_for_reference_proof(&self, _bucket: &str) -> Result> { + Ok(None) + } + } + #[tokio::test] async fn tier_config_update_path_acquires_meta_namespace_sidecar_lock_before_save() { let manager = TierConfigMgr::new(); @@ -11296,6 +11377,7 @@ mod tests { lock_manager: Arc, lock_requests: Mutex>, listed_versions: Mutex>, + lifecycle_configs: Mutex>, } impl Default for CasConfigStore { @@ -11316,6 +11398,7 @@ mod tests { lock_manager: Arc::new(rustfs_lock::GlobalLockManager::new()), lock_requests: Mutex::new(Vec::new()), listed_versions: Mutex::new(Vec::new()), + lifecycle_configs: Mutex::new(HashMap::new()), } } } @@ -11342,6 +11425,13 @@ mod tests { .push(object); } + fn add_lifecycle_config(&self, bucket: &str, config: BucketLifecycleConfiguration) { + self.lifecycle_configs + .lock() + .expect("tier reference fixture should not poison") + .insert(bucket.to_string(), config); + } + async fn insert_config_object(&self, object: String, data: Vec) { self.objects .lock() @@ -11917,6 +12007,18 @@ mod tests { } } + #[async_trait::async_trait] + impl TierReferenceProofStore for CasConfigStore { + async fn lifecycle_config_for_reference_proof(&self, bucket: &str) -> Result> { + Ok(self + .lifecycle_configs + .lock() + .expect("tier reference fixture should not poison") + .get(bucket) + .cloned()) + } + } + #[tokio::test] async fn remove_and_clear_full_update_paths_preserve_force() { let remove_store = Arc::new(CasConfigStore::default()); @@ -12090,6 +12192,96 @@ mod tests { assert!(err.message.contains("photos/2026/old-destination.jpg"), "{}", err.message); } + #[tokio::test] + async fn zero_reference_proof_blocks_lifecycle_transition_references() { + let current = build_rustfs_tier("COLD-A"); + let current_identity = tier_backend_identity(¤t).expect("current identity should encode"); + let target = TierMutationIntentTarget { + tier_name: "COLD-A".to_string(), + old_backend_identity: Some(current_identity), + new_backend_identity: None, + }; + let config = BucketLifecycleConfiguration { + expiry_updated_at: None, + rules: vec![LifecycleRule { + status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), + expiration: None, + abort_incomplete_multipart_upload: None, + del_marker_expiration: None, + filter: None, + id: Some("move-current".to_string()), + noncurrent_version_expiration: None, + noncurrent_version_transitions: None, + prefix: None, + transitions: Some(vec![Transition { + days: Some(1), + date: None, + storage_class: Some(TransitionStorageClass::from_static("COLD-A")), + }]), + }], + }; + let store = Arc::new(CasConfigStore::default()); + store.add_listed_version(ObjectInfo { + bucket: "photos".to_string(), + name: "safe.txt".to_string(), + ..Default::default() + }); + store.add_lifecycle_config("photos", config); + + let err = ensure_no_authoritative_tier_object_references(store, &[target]) + .await + .expect_err("lifecycle rule should block deletion"); + + assert_eq!(err.code, ERR_TIER_BACKEND_IN_USE.code); + assert!(err.message.contains("move-current"), "{}", err.message); + assert!(err.message.contains("photos"), "{}", err.message); + } + + #[tokio::test] + async fn zero_reference_proof_blocks_lifecycle_noncurrent_transition_references() { + let current = build_rustfs_tier("COLD-A"); + let current_identity = tier_backend_identity(¤t).expect("current identity should encode"); + let target = TierMutationIntentTarget { + tier_name: "COLD-A".to_string(), + old_backend_identity: Some(current_identity), + new_backend_identity: None, + }; + let config = BucketLifecycleConfiguration { + expiry_updated_at: None, + rules: vec![LifecycleRule { + status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), + expiration: None, + abort_incomplete_multipart_upload: None, + del_marker_expiration: None, + filter: None, + id: Some("move-noncurrent".to_string()), + noncurrent_version_expiration: None, + noncurrent_version_transitions: Some(vec![NoncurrentVersionTransition { + noncurrent_days: Some(1), + newer_noncurrent_versions: None, + storage_class: Some(TransitionStorageClass::from_static("COLD-A")), + }]), + prefix: None, + transitions: None, + }], + }; + let store = Arc::new(CasConfigStore::default()); + store.add_listed_version(ObjectInfo { + bucket: "photos".to_string(), + name: "safe.txt".to_string(), + ..Default::default() + }); + store.add_lifecycle_config("photos", config); + + let err = ensure_no_authoritative_tier_object_references(store, &[target]) + .await + .expect_err("lifecycle rule should block deletion"); + + assert_eq!(err.code, ERR_TIER_BACKEND_IN_USE.code); + assert!(err.message.contains("move-noncurrent"), "{}", err.message); + assert!(err.message.contains("photos"), "{}", err.message); + } + #[tokio::test] async fn zero_reference_proof_blocks_persisted_journal_transaction_and_free_version_references() { let current = build_rustfs_tier("COLD-A"); diff --git a/crates/s3-client/src/api_s3_datatypes.rs b/crates/s3-client/src/api_s3_datatypes.rs index d7797e96b..7e0657236 100644 --- a/crates/s3-client/src/api_s3_datatypes.rs +++ b/crates/s3-client/src/api_s3_datatypes.rs @@ -396,3 +396,30 @@ impl DeleteMultiObjects { }) } } + +#[cfg(test)] +mod tests { + use super::ListBucketV2Result; + + #[test] + fn list_bucket_v2_common_prefix_accepts_s3_pascal_case_xml() { + let xml = r#" + + tier-bucket + + 1 + 1000 + / + false + + tenant-a/ + + + "#; + + let result = quick_xml::de::from_str::(xml).expect("S3 list response should decode"); + + assert_eq!(result.common_prefixes.len(), 1); + assert_eq!(result.common_prefixes[0].prefix, "tenant-a/"); + } +} diff --git a/rustfs/src/app/object/get.rs b/rustfs/src/app/object/get.rs index 34d8f7f14..e02c6528e 100644 --- a/rustfs/src/app/object/get.rs +++ b/rustfs/src/app/object/get.rs @@ -3820,7 +3820,15 @@ impl DefaultObjectUsecase { Box::pin(self.execute_get_object_inner(req)) } + fn complete_get_object_error(helper: OperationHelper, err: S3Error) -> S3Result> { + let result = Err(err); + let _ = helper.complete(&result); + result + } + async fn execute_get_object_inner(&self, req: S3Request) -> S3Result> { + let helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject).suppress_event(); + if let Some(context) = &self.context { let _ = context.object_store(); } @@ -3840,7 +3848,10 @@ impl DefaultObjectUsecase { context.start_time.elapsed().as_secs_f64(), ); } - let bootstrap = self.init_get_object_bootstrap(&req.input.bucket, &req.input.key, &request_id)?; + let bootstrap = match self.init_get_object_bootstrap(&req.input.bucket, &req.input.key, &request_id) { + Ok(bootstrap) => bootstrap, + Err(err) => return Self::complete_get_object_error(helper, err), + }; record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_REQUEST_SHAPE, request_shape_start); let timeout_config = bootstrap.timeout_config; let wrapper = bootstrap.wrapper; @@ -3848,7 +3859,6 @@ impl DefaultObjectUsecase { let concurrent_requests = bootstrap.concurrent_requests; let mut lifecycle = GetObjectBodyLifecycle::tracked(bootstrap.request_guard); - let helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject).suppress_event(); // mc get 3 // Cheap request-shape validations run first so invalid requests keep @@ -3858,7 +3868,7 @@ impl DefaultObjectUsecase { Ok(validated) => validated, Err(err) => { lifecycle.finish_err(); - return Err(err); + return Self::complete_get_object_error(helper, err); } }; record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_REQUEST_VALIDATION, request_validation_start); @@ -3875,7 +3885,10 @@ impl DefaultObjectUsecase { let store_lookup_start = stage_metrics_enabled.then(std::time::Instant::now); let Some(store) = self.object_store() else { lifecycle.finish_err(); - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + return Self::complete_get_object_error( + helper, + S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()), + ); }; if let Some(store_lookup_start) = store_lookup_start { rustfs_io_metrics::record_get_object_stage_duration( @@ -3887,7 +3900,7 @@ impl DefaultObjectUsecase { let bucket_validation_start = stage_metrics_enabled.then(std::time::Instant::now); if let Err(err) = validate_bucket_exists(&store, &req.input.bucket).await { lifecycle.finish_err(); - return Err(err); + return Self::complete_get_object_error(helper, err); } record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_BUCKET_VALIDATION, bucket_validation_start); @@ -3896,7 +3909,7 @@ impl DefaultObjectUsecase { Ok(request_context) => request_context, Err(err) => { lifecycle.finish_err(); - return Err(err); + return Self::complete_get_object_error(helper, err); } }; if let Some(request_context_start) = request_context_start { @@ -3951,7 +3964,7 @@ impl DefaultObjectUsecase { return result; } lifecycle.finish_err(); - return Err(err); + return Self::complete_get_object_error(helper.version_id(version_id_for_event), err); } }; let GetObjectPreparedRead { io_planning, read_setup } = prepared_read; @@ -4040,7 +4053,7 @@ impl DefaultObjectUsecase { .await; let output_context = match output_context { Ok(output_context) => output_context, - Err(err) => return Err(err), + Err(err) => return Self::complete_get_object_error(helper.version_id(version_id_for_event), err), }; if let Some(output_build_start) = output_build_start { rustfs_io_metrics::record_get_object_stage_duration( diff --git a/rustfs/src/storage/helper.rs b/rustfs/src/storage/helper.rs index 5b6a2f88e..b8eb3c15d 100644 --- a/rustfs/src/storage/helper.rs +++ b/rustfs/src/storage/helper.rs @@ -23,7 +23,7 @@ use http::StatusCode; use metrics::counter; use rustfs_audit::{ ObjectVersion, - entity::{ApiDetails, ApiDetailsBuilder, AuditEntryBuilder}, + entity::{ApiDetailsBuilder, AuditEntryBuilder}, global::AuditLogger, }; use rustfs_io_metrics::record_s3_op; @@ -136,7 +136,7 @@ impl OperationHelper { record_s3_op(op); - // Fast path: when both chains are disabled, avoid all request parsing/builder work. + // Fast path: when both chains are disabled, avoid audit/notify builder work. if !audit_enabled && !notify_enabled { return Self::Disabled; } @@ -180,7 +180,7 @@ impl OperationHelper { let audit_builder = if audit_enabled { Some( - AuditEntryBuilder::new("1.0", event, trigger, ApiDetails::default()) + AuditEntryBuilder::new("1.0", event, trigger, api_builder.clone().build()) .remote_host(remote_host) .user_agent(get_request_user_agent(&req.headers)) .req_host(get_request_host(&req.headers)) @@ -453,8 +453,8 @@ mod tests { use rustfs_s3_ops::S3Operation; use rustfs_s3_types::EventName; use rustfs_utils::http::headers::{AMZ_REQUEST_ID, REQUEST_ID_HEADER}; - use s3s::dto::{DeleteObjectTaggingInput, DeleteObjectTaggingOutput}; - use s3s::{S3Request, S3Response}; + use s3s::dto::{DeleteObjectTaggingInput, DeleteObjectTaggingOutput, GetObjectInput, GetObjectOutput}; + use s3s::{S3Error, S3ErrorCode, S3Request, S3Response}; use std::sync::{Arc, Mutex}; use temp_env::{async_with_vars, with_vars}; @@ -615,6 +615,71 @@ mod tests { ); } + #[test] + fn operation_helper_initializes_audit_api_details_before_completion() { + with_vars( + [ + (rustfs_config::ENV_NOTIFY_ENABLE, Some("false")), + (rustfs_config::ENV_AUDIT_ENABLE, Some("true")), + ], + || { + refresh_notify_module_enabled(); + refresh_audit_module_enabled(); + + let input = GetObjectInput::builder() + .bucket("audit-bucket".to_string()) + .key("missing/object.txt".to_string()) + .build() + .expect("get object input should build"); + let req = build_request(input, Method::GET, Uri::from_static("/audit-bucket/missing/object.txt")); + let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject); + + let OperationHelper::Enabled(state) = &mut helper else { + panic!("helper should be enabled when the audit switch is on"); + }; + let audit_entry = state.audit_builder.take().expect("audit builder should exist").build(); + + assert_eq!(audit_entry.api.name.as_deref(), Some("s3:GetObject")); + assert_eq!(audit_entry.api.bucket.as_deref(), Some("audit-bucket")); + assert_eq!(audit_entry.api.object.as_deref(), Some("missing/object.txt")); + assert!(audit_entry.api.status_code.is_none()); + }, + ); + } + + #[test] + fn operation_helper_complete_records_failed_status_code() { + with_vars( + [ + (rustfs_config::ENV_NOTIFY_ENABLE, Some("false")), + (rustfs_config::ENV_AUDIT_ENABLE, Some("true")), + ], + || { + refresh_notify_module_enabled(); + refresh_audit_module_enabled(); + + let input = GetObjectInput::builder() + .bucket("audit-bucket".to_string()) + .key("missing/object.txt".to_string()) + .build() + .expect("get object input should build"); + let req = build_request(input, Method::GET, Uri::from_static("/audit-bucket/missing/object.txt")); + let result: Result, S3Error> = Err(S3Error::new(S3ErrorCode::NoSuchKey)); + let mut helper = + OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject).complete(&result); + + let OperationHelper::Enabled(state) = &mut helper else { + panic!("helper should be enabled when the audit switch is on"); + }; + let audit_entry = state.audit_builder.take().expect("audit builder should exist").build(); + + assert_eq!(audit_entry.api.name.as_deref(), Some("s3:GetObject")); + assert_eq!(audit_entry.api.status.as_deref(), Some("failure")); + assert_eq!(audit_entry.api.status_code, Some(404)); + }, + ); + } + #[test] fn operation_helper_prioritizes_request_context_for_request_id() { with_vars(