diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index b7ce10f1b..23639060d 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -45,14 +45,17 @@ use rustfs_ecstore::bucket::{ bucket_lifecycle_ops::{RestoreRequestOps, post_restore_opts}, lifecycle::{self, TransitionOptions}, }, + metadata::{BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG}, metadata_sys, - object_lock::objectlock_sys::{BucketObjectLockSys, check_object_lock_for_deletion}, + object_lock::objectlock_sys::{BucketObjectLockSys, check_object_lock_for_deletion, check_retention_for_modification}, quota::QuotaOperation, replication::{ DeletedObjectReplicationInfo, get_must_replicate_options, must_replicate, schedule_replication, schedule_replication_delete, }, tagging::{decode_tags, encode_tags}, + utils::serialize, + versioning::VersioningApi, versioning_sys::BucketVersioningSys, }; use rustfs_ecstore::client::object_api_utils::to_s3s_etag; @@ -60,7 +63,7 @@ use rustfs_ecstore::compress::{MIN_COMPRESSIBLE_SIZE, is_compressible}; use rustfs_ecstore::error::{StorageError, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found}; use rustfs_ecstore::new_object_layer_fn; use rustfs_ecstore::set_disk::is_valid_storage_class; -use rustfs_ecstore::store_api::{HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOptions, PutObjReader}; +use rustfs_ecstore::store_api::{BucketOptions, HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOptions, PutObjReader}; use rustfs_filemeta::{ REPLICATE_INCOMING_DELETE, ReplicationStatusType, ReplicationType, RestoreStatusOps, VersionPurgeStatusType, parse_restore_obj_status, @@ -721,6 +724,220 @@ impl DefaultObjectUsecase { Ok(S3Response::new(PutObjectAclOutput::default())) } + pub async fn execute_put_object_legal_hold( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutLegalHold, "s3:PutObjectLegalHold"); + let PutObjectLegalHoldInput { + bucket, + key, + legal_hold, + version_id, + .. + } = req.input.clone(); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let _ = store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + validate_bucket_object_lock_enabled(&bucket).await?; + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) + .await + .map_err(ApiError::from)?; + + let eval_metadata = parse_object_lock_legal_hold(legal_hold)?; + + let popts = ObjectOptions { + mod_time: opts.mod_time, + version_id: opts.version_id, + eval_metadata: Some(eval_metadata), + ..Default::default() + }; + + let info = store.put_object_metadata(&bucket, &key, &popts).await.map_err(|e| { + error!("put_object_metadata failed, {}", e.to_string()); + s3_error!(InternalError, "{}", e.to_string()) + })?; + + let output = PutObjectLegalHoldOutput { + request_charged: Some(RequestCharged::from_static(RequestCharged::REQUESTER)), + }; + let version_id = req.input.version_id.clone().unwrap_or_default(); + helper = helper.object(info).version_id(version_id); + + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result + } + + #[instrument(level = "debug", skip(self))] + pub async fn execute_put_object_lock_configuration( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let PutObjectLockConfigurationInput { + bucket, + object_lock_configuration, + .. + } = req.input; + + let Some(input_cfg) = object_lock_configuration else { return Err(s3_error!(InvalidArgument)) }; + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + match metadata_sys::get_object_lock_config(&bucket).await { + Ok(_) => {} + Err(err) => { + if err == StorageError::ConfigNotFound { + return Err(S3Error::with_message( + S3ErrorCode::InvalidBucketState, + "Object Lock configuration cannot be enabled on existing buckets".to_string(), + )); + } + warn!("get_object_lock_config err {:?}", err); + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + "Failed to get bucket ObjectLockConfiguration".to_string(), + )); + } + }; + + let data = serialize(&input_cfg).map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("{}", err)))?; + + metadata_sys::update(&bucket, OBJECT_LOCK_CONFIG, data) + .await + .map_err(ApiError::from)?; + + // When Object Lock is enabled, automatically enable versioning if not already enabled. + // This matches AWS S3 and MinIO behavior. + let versioning_config = BucketVersioningSys::get(&bucket).await.map_err(ApiError::from)?; + if !versioning_config.enabled() { + let enable_versioning_config = VersioningConfiguration { + status: Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)), + ..Default::default() + }; + let versioning_data = serialize(&enable_versioning_config) + .map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("{}", err)))?; + metadata_sys::update(&bucket, BUCKET_VERSIONING_CONFIG, versioning_data) + .await + .map_err(ApiError::from)?; + } + + Ok(S3Response::new(PutObjectLockConfigurationOutput::default())) + } + + pub async fn execute_put_object_retention( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutRetention, "s3:PutObjectRetention"); + let PutObjectRetentionInput { + bucket, + key, + retention, + version_id, + .. + } = req.input.clone(); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + validate_bucket_object_lock_enabled(&bucket).await?; + + let new_retain_until = retention + .as_ref() + .and_then(|r| r.retain_until_date.as_ref()) + .map(|d| OffsetDateTime::from(d.clone())); + + // TODO(security): Known TOCTOU race condition (fix in future PR). + // + // There is a time-of-check-time-of-use (TOCTOU) window between the retention + // check below (using get_object_info + check_retention_for_modification) and + // the actual update performed later in put_object_metadata. + // + // In theory: + // * Thread A reads retention mode = GOVERNANCE and checks the bypass header. + // * Thread B updates retention to COMPLIANCE mode. + // * Thread A then proceeds to modify retention, still assuming GOVERNANCE, + // and effectively bypasses what is now COMPLIANCE mode. + // + // This would violate the S3 spec, which states that COMPLIANCE-mode retention + // cannot be modified even with a bypass header. + // + // Possible fixes (to be implemented in a future change): + // 1. Pass the expected retention mode down to the storage layer and verify + // it has not changed immediately before the update. + // 2. Use optimistic concurrency (e.g., version/etag) so that the update + // fails if the object changed between check and update. + // 3. Perform the retention check inside the same lock/transaction scope as + // the metadata update within the storage layer. + // + // Current mitigation: the storage layer provides a fast_lock_manager, which + // offers some protection, but it does not fully eliminate this race. + let check_opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) + .await + .map_err(ApiError::from)?; + + if let Ok(existing_obj_info) = store.get_object_info(&bucket, &key, &check_opts).await { + let bypass_governance = has_bypass_governance_header(&req.headers); + if let Some(block_reason) = + check_retention_for_modification(&existing_obj_info.user_defined, new_retain_until, bypass_governance) + { + return Err(S3Error::with_message(S3ErrorCode::AccessDenied, block_reason.error_message())); + } + } + + let eval_metadata = parse_object_lock_retention(retention)?; + + let mut opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) + .await + .map_err(ApiError::from)?; + opts.eval_metadata = Some(eval_metadata); + + let object_info = store.put_object_metadata(&bucket, &key, &opts).await.map_err(|e| { + error!("put_object_metadata failed, {}", e.to_string()); + s3_error!(InternalError, "{}", e.to_string()) + })?; + + let output = PutObjectRetentionOutput { + request_charged: Some(RequestCharged::from_static(RequestCharged::REQUESTER)), + }; + + let version_id = req.input.version_id.clone().unwrap_or_else(|| Uuid::new_v4().to_string()); + helper = helper.object(object_info).version_id(version_id); + + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result + } + #[instrument(level = "debug", skip(self, req))] pub async fn execute_put_object_tagging( &self, @@ -1551,6 +1768,189 @@ impl DefaultObjectUsecase { })) } + pub async fn execute_get_object_attributes( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedAttributes, "s3:GetObjectAttributes"); + let GetObjectAttributesInput { bucket, key, .. } = req.input.clone(); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + if let Err(e) = store + .get_object_reader(&bucket, &key, None, HeaderMap::new(), &ObjectOptions::default()) + .await + { + return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{e}"))); + } + + let output = GetObjectAttributesOutput { + delete_marker: None, + object_parts: None, + ..Default::default() + }; + + let version_id = req.input.version_id.clone().unwrap_or_default(); + helper = helper + .object(ObjectInfo { + name: key.clone(), + bucket, + ..Default::default() + }) + .version_id(version_id); + + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result + } + + pub async fn execute_get_object_legal_hold( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGetLegalHold, "s3:GetObjectLegalHold"); + let GetObjectLegalHoldInput { + bucket, key, version_id, .. + } = req.input.clone(); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let _ = store + .get_bucket_info(&bucket, &BucketOptions::default()) + .await + .map_err(ApiError::from)?; + + validate_bucket_object_lock_enabled(&bucket).await?; + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) + .await + .map_err(ApiError::from)?; + + let object_info = store.get_object_info(&bucket, &key, &opts).await.map_err(|e| { + error!("get_object_info failed, {}", e.to_string()); + s3_error!(InternalError, "{}", e.to_string()) + })?; + + let legal_hold = object_info + .user_defined + .get(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER) + .map(|v| v.as_str().to_string()); + + let status = if let Some(v) = legal_hold { + v + } else { + ObjectLockLegalHoldStatus::OFF.to_string() + }; + + let output = GetObjectLegalHoldOutput { + legal_hold: Some(ObjectLockLegalHold { + status: Some(ObjectLockLegalHoldStatus::from(status)), + }), + }; + + let version_id = req.input.version_id.clone().unwrap_or_else(|| Uuid::new_v4().to_string()); + helper = helper.object(object_info).version_id(version_id); + + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result + } + + #[instrument(level = "debug", skip(self))] + pub async fn execute_get_object_lock_configuration( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let GetObjectLockConfigurationInput { bucket, .. } = req.input; + + let object_lock_configuration = match metadata_sys::get_object_lock_config(&bucket).await { + Ok((cfg, _created)) => Some(cfg), + Err(err) => { + if err == StorageError::ConfigNotFound { + return Err(S3Error::with_message( + S3ErrorCode::ObjectLockConfigurationNotFoundError, + "Object Lock configuration does not exist for this bucket".to_string(), + )); + } + warn!("get_object_lock_config err {:?}", err); + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + "Failed to load Object Lock configuration".to_string(), + )); + } + }; + + Ok(S3Response::new(GetObjectLockConfigurationOutput { + object_lock_configuration, + })) + } + + pub async fn execute_get_object_retention( + &self, + req: S3Request, + ) -> S3Result> { + if let Some(context) = &self.context { + let _ = context.object_store(); + } + + let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGetRetention, "s3:GetObjectRetention"); + let GetObjectRetentionInput { + bucket, key, version_id, .. + } = req.input.clone(); + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + validate_bucket_object_lock_enabled(&bucket).await?; + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) + .await + .map_err(ApiError::from)?; + + let object_info = store.get_object_info(&bucket, &key, &opts).await.map_err(|e| { + error!("get_object_info failed, {}", e.to_string()); + s3_error!(InternalError, "{}", e.to_string()) + })?; + + let mode = object_info + .user_defined + .get("x-amz-object-lock-mode") + .map(|v| ObjectLockRetentionMode::from(v.as_str().to_string())); + + let retain_until_date = object_info + .user_defined + .get("x-amz-object-lock-retain-until-date") + .and_then(|v| OffsetDateTime::parse(v.as_str(), &Rfc3339).ok()) + .map(Timestamp::from); + + let output = GetObjectRetentionOutput { + retention: Some(ObjectLockRetention { mode, retain_until_date }), + }; + let version_id = req.input.version_id.clone().unwrap_or_default(); + helper = helper.object(object_info).version_id(version_id); + + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result + } + #[instrument(level = "debug", skip(self, req))] pub async fn execute_get_object_tagging( &self, @@ -2895,6 +3295,51 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[tokio::test] + async fn execute_get_object_attributes_returns_internal_error_when_store_uninitialized() { + let input = GetObjectAttributesInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::GET); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_get_object_attributes(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_get_object_legal_hold_returns_internal_error_when_store_uninitialized() { + let input = GetObjectLegalHoldInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::GET); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_get_object_legal_hold(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_get_object_retention_returns_internal_error_when_store_uninitialized() { + let input = GetObjectRetentionInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::GET); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_get_object_retention(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + #[tokio::test] async fn execute_get_object_tagging_returns_internal_error_when_store_uninitialized() { let input = GetObjectTaggingInput::builder() @@ -2925,6 +3370,54 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::InternalError); } + #[tokio::test] + async fn execute_put_object_legal_hold_returns_internal_error_when_store_uninitialized() { + let input = PutObjectLegalHoldInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::PUT); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_put_object_legal_hold(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_put_object_lock_configuration_returns_internal_error_when_store_uninitialized() { + let input = PutObjectLockConfigurationInput::builder() + .bucket("test-bucket".to_string()) + .object_lock_configuration(Some(ObjectLockConfiguration { + object_lock_enabled: Some(ObjectLockEnabled::from_static(ObjectLockEnabled::ENABLED)), + rule: None, + })) + .build() + .unwrap(); + + let req = build_request(input, Method::PUT); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_put_object_lock_configuration(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[tokio::test] + async fn execute_put_object_retention_returns_internal_error_when_store_uninitialized() { + let input = PutObjectRetentionInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + + let req = build_request(input, Method::PUT); + let usecase = DefaultObjectUsecase::without_context(); + + let err = usecase.execute_put_object_retention(req).await.unwrap_err(); + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + #[tokio::test] async fn execute_put_object_tagging_returns_internal_error_when_store_uninitialized() { let input = PutObjectTaggingInput::builder() diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 14626e05d..9d0aca8de 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -21,21 +21,19 @@ use crate::storage::helper::OperationHelper; use crate::storage::options::get_content_sha256; use crate::storage::{ access::{ReqInfo, authorize_request, has_bypass_governance_header}, - options::{copy_src_opts, del_opts, extract_metadata, get_opts, parse_copy_source_range}, + options::{copy_src_opts, del_opts, extract_metadata, parse_copy_source_range}, }; use crate::storage::{ decrypt_managed_encryption_key, derive_part_nonce, get_buffer_size_opt_in, get_validated_store, has_replication_rules, - parse_object_lock_legal_hold, parse_object_lock_retention, validate_bucket_object_lock_enabled, }; use futures::StreamExt; use http::HeaderMap; use rustfs_ecstore::{ bucket::{ - metadata::{BUCKET_ACL_CONFIG, BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG}, + metadata::BUCKET_ACL_CONFIG, metadata_sys, - object_lock::objectlock_sys::{check_object_lock_for_deletion, check_retention_for_modification}, + object_lock::objectlock_sys::check_object_lock_for_deletion, replication::{DeletedObjectReplicationInfo, check_replicate_delete, schedule_replication_delete}, - utils::serialize, versioning::VersioningApi, versioning_sys::BucketVersioningSys, }, @@ -66,34 +64,19 @@ use rustfs_targets::EventName; use rustfs_utils::{ CompressionAlgorithm, extract_params_header, extract_resp_elements, get_request_host, get_request_port, get_request_user_agent, - http::{ - AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER, - headers::{AMZ_DECODED_CONTENT_LENGTH, RESERVED_METADATA_PREFIX_LOWER}, - }, + http::headers::{AMZ_DECODED_CONTENT_LENGTH, RESERVED_METADATA_PREFIX_LOWER}, path::is_dir_object, }; use rustfs_zip::CompressionFormat; use s3s::{S3, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, dto::*, s3_error}; use serde::{Deserialize, Serialize}; use std::{collections::HashMap, fmt::Debug, path::Path, sync::LazyLock}; -use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tokio::io::{AsyncRead, AsyncSeek}; use tokio_tar::Archive; use tokio_util::io::StreamReader; use tracing::{debug, error, instrument, warn}; use uuid::Uuid; -macro_rules! try_ { - ($result:expr) => { - match $result { - Ok(val) => val, - Err(err) => { - return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{}", err))); - } - } - }; -} - const DEFAULT_OWNER_ID: &str = "rustfsadmin"; const DEFAULT_OWNER_DISPLAY_NAME: &str = "RustFS Tester"; const DEFAULT_OWNER_EMAIL: &str = "tester@rustfs.local"; @@ -1547,93 +1530,16 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedAttributes, "s3:GetObjectAttributes"); - let GetObjectAttributesInput { bucket, key, .. } = req.input.clone(); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - if let Err(e) = store - .get_object_reader(&bucket, &key, None, HeaderMap::new(), &ObjectOptions::default()) - .await - { - return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("{e}"))); - } - - let output = GetObjectAttributesOutput { - delete_marker: None, - object_parts: None, - ..Default::default() - }; - - let version_id = req.input.version_id.clone().unwrap_or_default(); - helper = helper - .object(ObjectInfo { - name: key.clone(), - bucket, - ..Default::default() - }) - .version_id(version_id); - - let result = Ok(S3Response::new(output)); - let _ = helper.complete(&result); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_get_object_attributes(req).await } async fn get_object_legal_hold( &self, req: S3Request, ) -> S3Result> { - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGetLegalHold, "s3:GetObjectLegalHold"); - let GetObjectLegalHoldInput { - bucket, key, version_id, .. - } = req.input.clone(); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let _ = store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - // check object lock - validate_bucket_object_lock_enabled(&bucket).await?; - - let opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) - .await - .map_err(ApiError::from)?; - - let object_info = store.get_object_info(&bucket, &key, &opts).await.map_err(|e| { - error!("get_object_info failed, {}", e.to_string()); - s3_error!(InternalError, "{}", e.to_string()) - })?; - - let legal_hold = object_info - .user_defined - .get(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER) - .map(|v| v.as_str().to_string()); - - let status = if let Some(v) = legal_hold { - v - } else { - ObjectLockLegalHoldStatus::OFF.to_string() - }; - - let output = GetObjectLegalHoldOutput { - legal_hold: Some(ObjectLockLegalHold { - status: Some(ObjectLockLegalHoldStatus::from(status)), - }), - }; - - let version_id = req.input.version_id.clone().unwrap_or_else(|| Uuid::new_v4().to_string()); - helper = helper.object(object_info).version_id(version_id); - - let result = Ok(S3Response::new(output)); - let _ = helper.complete(&result); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_get_object_legal_hold(req).await } #[instrument(level = "debug", skip(self))] @@ -1641,77 +1547,16 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let GetObjectLockConfigurationInput { bucket, .. } = req.input; - - let object_lock_configuration = match metadata_sys::get_object_lock_config(&bucket).await { - Ok((cfg, _created)) => Some(cfg), - Err(err) => { - if err == StorageError::ConfigNotFound { - return Err(S3Error::with_message( - S3ErrorCode::ObjectLockConfigurationNotFoundError, - "Object Lock configuration does not exist for this bucket".to_string(), - )); - } - warn!("get_object_lock_config err {:?}", err); - return Err(S3Error::with_message( - S3ErrorCode::InternalError, - "Failed to load Object Lock configuration".to_string(), - )); - } - }; - - // warn!("object_lock_configuration {:?}", &object_lock_configuration); - - Ok(S3Response::new(GetObjectLockConfigurationOutput { - object_lock_configuration, - })) + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_get_object_lock_configuration(req).await } async fn get_object_retention( &self, req: S3Request, ) -> S3Result> { - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGetRetention, "s3:GetObjectRetention"); - let GetObjectRetentionInput { - bucket, key, version_id, .. - } = req.input.clone(); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - // check object lock - validate_bucket_object_lock_enabled(&bucket).await?; - - let opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) - .await - .map_err(ApiError::from)?; - - let object_info = store.get_object_info(&bucket, &key, &opts).await.map_err(|e| { - error!("get_object_info failed, {}", e.to_string()); - s3_error!(InternalError, "{}", e.to_string()) - })?; - - let mode = object_info - .user_defined - .get("x-amz-object-lock-mode") - .map(|v| ObjectLockRetentionMode::from(v.as_str().to_string())); - - let retain_until_date = object_info - .user_defined - .get("x-amz-object-lock-retain-until-date") - .and_then(|v| OffsetDateTime::parse(v.as_str(), &Rfc3339).ok()) - .map(Timestamp::from); - - let output = GetObjectRetentionOutput { - retention: Some(ObjectLockRetention { mode, retain_until_date }), - }; - let version_id = req.input.version_id.clone().unwrap_or_default(); - helper = helper.object(object_info).version_id(version_id); - - let result = Ok(S3Response::new(output)); - let _ = helper.complete(&result); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_get_object_retention(req).await } #[instrument(level = "debug", skip(self))] @@ -2244,53 +2089,8 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutLegalHold, "s3:PutObjectLegalHold"); - let PutObjectLegalHoldInput { - bucket, - key, - legal_hold, - version_id, - .. - } = req.input.clone(); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let _ = store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - // check object lock - validate_bucket_object_lock_enabled(&bucket).await?; - let opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) - .await - .map_err(ApiError::from)?; - - let eval_metadata = parse_object_lock_legal_hold(legal_hold)?; - - let popts = ObjectOptions { - mod_time: opts.mod_time, - version_id: opts.version_id, - eval_metadata: Some(eval_metadata), - ..Default::default() - }; - - let info = store.put_object_metadata(&bucket, &key, &popts).await.map_err(|e| { - error!("put_object_metadata failed, {}", e.to_string()); - s3_error!(InternalError, "{}", e.to_string()) - })?; - - let output = PutObjectLegalHoldOutput { - request_charged: Some(RequestCharged::from_static(RequestCharged::REQUESTER)), - }; - let version_id = req.input.version_id.clone().unwrap_or_default(); - helper = helper.object(info).version_id(version_id); - - let result = Ok(S3Response::new(output)); - let _ = helper.complete(&result); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_put_object_legal_hold(req).await } #[instrument(level = "debug", skip(self))] @@ -2298,141 +2098,16 @@ impl S3 for FS { &self, req: S3Request, ) -> S3Result> { - let PutObjectLockConfigurationInput { - bucket, - object_lock_configuration, - .. - } = req.input; - - let Some(input_cfg) = object_lock_configuration else { return Err(s3_error!(InvalidArgument)) }; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - store - .get_bucket_info(&bucket, &BucketOptions::default()) - .await - .map_err(ApiError::from)?; - - let _ = match metadata_sys::get_object_lock_config(&bucket).await { - Ok(_) => {} - Err(err) => { - if err == StorageError::ConfigNotFound { - return Err(S3Error::with_message( - S3ErrorCode::InvalidBucketState, - "Object Lock configuration cannot be enabled on existing buckets".to_string(), - )); - } - warn!("get_object_lock_config err {:?}", err); - return Err(S3Error::with_message( - S3ErrorCode::InternalError, - "Failed to get bucket ObjectLockConfiguration".to_string(), - )); - } - }; - - let data = try_!(serialize(&input_cfg)); - - metadata_sys::update(&bucket, OBJECT_LOCK_CONFIG, data) - .await - .map_err(ApiError::from)?; - - // When Object Lock is enabled, automatically enable versioning if not already enabled - // This matches AWS S3 and MinIO behavior - let versioning_config = BucketVersioningSys::get(&bucket).await.map_err(ApiError::from)?; - if !versioning_config.enabled() { - let enable_versioning_config = VersioningConfiguration { - status: Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)), - ..Default::default() - }; - let versioning_data = try_!(serialize(&enable_versioning_config)); - metadata_sys::update(&bucket, BUCKET_VERSIONING_CONFIG, versioning_data) - .await - .map_err(ApiError::from)?; - } - - Ok(S3Response::new(PutObjectLockConfigurationOutput::default())) + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_put_object_lock_configuration(req).await } async fn put_object_retention( &self, req: S3Request, ) -> S3Result> { - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutRetention, "s3:PutObjectRetention"); - let PutObjectRetentionInput { - bucket, - key, - retention, - version_id, - .. - } = req.input.clone(); - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - // check object lock - validate_bucket_object_lock_enabled(&bucket).await?; - - // Extract new retain_until_date for validation - let new_retain_until = retention - .as_ref() - .and_then(|r| r.retain_until_date.as_ref()) - .map(|d| OffsetDateTime::from(d.clone())); - - // Check if object already has retention and if modification is allowed - // This follows AWS S3 and MinIO behavior: - // - COMPLIANCE mode: can only extend retention period, never shorten - // - GOVERNANCE mode: requires bypass header to modify/shorten retention - // - // TODO: Known race condition (fix in future PR) - // There's a TOCTOU (time-of-check-time-of-use) window between retention check here - // and the actual update at put_object_metadata. In theory: - // Thread A: reads GOVERNANCE mode, checks bypass header - // Thread B: updates retention to COMPLIANCE mode - // Thread A: modifies retention, bypassing what is now COMPLIANCE mode - // This violates S3 spec that COMPLIANCE cannot be modified even with bypass. - // Fix options: - // 1. Pass expected retention mode to storage layer, verify before update - // 2. Use optimistic concurrency with version/etag checks - // 3. Perform check within the same lock scope as update in storage layer - // Current mitigation: Storage layer has fast_lock_manager which provides some protection - let check_opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) - .await - .map_err(ApiError::from)?; - - if let Ok(existing_obj_info) = store.get_object_info(&bucket, &key, &check_opts).await { - let bypass_governance = has_bypass_governance_header(&req.headers); - if let Some(block_reason) = - check_retention_for_modification(&existing_obj_info.user_defined, new_retain_until, bypass_governance) - { - return Err(S3Error::with_message(S3ErrorCode::AccessDenied, block_reason.error_message())); - } - } - - let eval_metadata = parse_object_lock_retention(retention)?; - - let mut opts: ObjectOptions = get_opts(&bucket, &key, version_id, None, &req.headers) - .await - .map_err(ApiError::from)?; - opts.eval_metadata = Some(eval_metadata); - - let object_info = store.put_object_metadata(&bucket, &key, &opts).await.map_err(|e| { - error!("put_object_metadata failed, {}", e.to_string()); - s3_error!(InternalError, "{}", e.to_string()) - })?; - - let output = PutObjectRetentionOutput { - request_charged: Some(RequestCharged::from_static(RequestCharged::REQUESTER)), - }; - - let version_id = req.input.version_id.clone().unwrap_or_else(|| Uuid::new_v4().to_string()); - helper = helper.object(object_info).version_id(version_id); - - let result = Ok(S3Response::new(output)); - let _ = helper.complete(&result); - result + let usecase = DefaultObjectUsecase::from_global(); + usecase.execute_put_object_retention(req).await } #[instrument(level = "debug", skip(self, req))]