From ff77cb55407fec41bdb09f5aeca20c63787efa43 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Thu, 9 Apr 2026 12:54:01 +0800 Subject: [PATCH] refactor(app): simplify object usecase plumbing (#2438) --- rustfs/src/app/object_usecase.rs | 233 ++++++--- rustfs/src/app/object_usecase/app_adapters.rs | 461 ------------------ .../src/app/object_usecase/get_object_flow.rs | 408 +++++++++++----- .../object_usecase/get_object_zero_copy.rs | 89 ++-- .../app/object_usecase/put_object_extract.rs | 72 +-- .../src/app/object_usecase/put_object_flow.rs | 116 ++--- rustfs/src/app/object_usecase/types.rs | 8 - .../src/app/object_usecase/zero_copy_tests.rs | 202 ++++---- 8 files changed, 661 insertions(+), 928 deletions(-) delete mode 100644 rustfs/src/app/object_usecase/app_adapters.rs diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 71669260b..04d0180d8 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -14,7 +14,6 @@ //! Object application use-case contracts. -mod app_adapters; mod get_object_flow; mod get_object_zero_copy; mod put_object_extract; @@ -22,8 +21,7 @@ mod put_object_flow; mod types; #[cfg(test)] mod zero_copy_tests; -use self::app_adapters::*; -use self::get_object_flow::{GetObjectBootstrap, GetObjectFlowRuntime}; +use self::get_object_flow::GetObjectBootstrap; use self::types::*; use crate::app::context::{AppContext, default_notify_interface, get_global_app_context}; @@ -158,6 +156,26 @@ impl Drop for DeadlockRequestGuard { } } +async fn resolve_bucket_default_server_side_encryption(bucket: &str) -> (Option, Option) { + let Some((config, _timestamp)) = metadata_sys::get_sse_config(bucket).await.ok() else { + return (None, None); + }; + let Some(default_sse) = config + .rules + .first() + .and_then(|rule| rule.apply_server_side_encryption_by_default.as_ref()) + else { + return (None, None); + }; + + let server_side_encryption = Some(match default_sse.sse_algorithm.as_str() { + "aws:kms" => ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS), + _ => ServerSideEncryption::from_static(ServerSideEncryption::AES256), + }); + + (server_side_encryption, default_sse.kms_master_key_id.clone()) +} + async fn enqueue_transitioned_delete_cleanup(bucket: &str, object: &str, opts: &ObjectOptions, existing: Option<&ObjectInfo>) { let Some(existing) = existing else { return; @@ -267,10 +285,6 @@ mod deadlock_request_guard_tests { assert_eq!(detector.tracked_count(), 0); } } -async fn maybe_enqueue_transition_immediate(obj_info: &ObjectInfo, src: LcEventSrc) { - enqueue_transition_immediate(obj_info, src).await; -} - fn normalize_delete_objects_version_id(version_id: Option) -> Result<(Option, Option), String> { let version_id = version_id.map(|v| v.trim().to_string()).filter(|v| !v.is_empty()); match version_id { @@ -605,9 +619,24 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let request_context = prepare_put_object_request_context(&req); - let (event_name, quota_operation, request_method_name) = put_object_execution_context(&req); - let helper = new_operation_helper(&req, event_name, S3Operation::PutObject, false); + let request_context = PutObjectRequestContext { + headers: req.headers.clone(), + trailing_headers: req.trailing_headers.clone(), + uri_query: req.uri.query().map(str::to_string), + is_post_object: req.extensions.get::().is_some(), + method: req.method.clone(), + uri: req.uri.clone(), + extensions: req.extensions.clone(), + credentials: req.credentials.clone(), + region: req.region.clone(), + service: req.service.clone(), + }; + let (event_name, quota_operation) = if request_context.is_post_object { + (EventName::ObjectCreatedPost, QuotaOperation::PostObject) + } else { + (EventName::ObjectCreatedPut, QuotaOperation::PutObject) + }; + let helper = OperationHelper::new(&req, event_name, S3Operation::PutObject); if request_context.is_post_object && is_post_object_sse_kms_requested(&req.input, &request_context.headers) { return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for POST object uploads")); @@ -626,10 +655,17 @@ impl DefaultObjectUsecase { .await?; let input = req.input; - let flow_result = - DefaultObjectUsecase::run_put_object_flow(input, request_context, request_method_name, resolved_size).await?; - let helper = bind_helper_object(helper, flow_result.helper_object, flow_result.helper_version_id); - complete_put_response(helper, flow_result.output) + let (output, helper_object) = DefaultObjectUsecase::run_put_object_flow(input, request_context, resolved_size).await?; + let helper_version_id = helper_object.version_id.map(|version_id| version_id.to_string()); + let helper = helper.object(helper_object); + let helper = if let Some(version_id) = helper_version_id { + helper.version_id(version_id) + } else { + helper + }; + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result } pub async fn execute_put_object_acl(&self, req: S3Request) -> S3Result> { @@ -1023,17 +1059,106 @@ impl DefaultObjectUsecase { .get::() .map(|ctx| ctx.request_id.clone()) .unwrap_or_else(|| request_context::RequestContext::fallback().request_id); - let bootstrap = init_get_object_bootstrap(&req.input.bucket, &req.input.key, &request_id)?; - let request_context = prepare_get_object_request_context(&req).await?; + let bootstrap = { + let timeout_config = TimeoutConfig::from_env(); + let wrapper = RequestTimeoutWrapper::with_request_id(timeout_config.clone(), request_id.clone()); + let request_start = std::time::Instant::now(); + let request_guard = crate::storage::concurrency::ConcurrencyManager::track_request(); + let concurrent_requests = GetObjectGuard::concurrent_requests(); + + let deadlock_detector = deadlock_detector::get_deadlock_detector(); + deadlock_detector.register_request(&request_id, format!("GetObject {}/{}", req.input.bucket, req.input.key)); + let deadlock_request_guard = DeadlockRequestGuard::new(deadlock_detector, request_id); + + if wrapper.is_timeout() { + warn!( + bucket = %req.input.bucket, + key = %req.input.key, + timeout_secs = timeout_config.get_object_timeout.as_secs(), + elapsed_ms = wrapper.elapsed().as_millis(), + "GetObject request timed out before processing" + ); + return Err(s3_error!(InternalError, "Request timeout before processing")); + } + + rustfs_io_metrics::record_get_object_request_start(concurrent_requests); + + debug!( + "GetObject request started with {} concurrent requests, timeout={:?}", + concurrent_requests, timeout_config.get_object_timeout + ); + + GetObjectBootstrap { + timeout_config, + wrapper, + request_start, + request_guard, + _deadlock_request_guard: deadlock_request_guard, + } + }; + let (request_context, version_id_for_event) = { + let GetObjectInput { + bucket, + key, + version_id, + part_number, + range, + .. + } = req.input.clone(); + + validate_object_key(&key, "GET")?; + + let part_number = part_number.map(|value| value as usize); + if let Some(part_number) = part_number + && part_number == 0 + { + return Err(s3_error!(InvalidArgument, "Invalid part number: part number must be greater than 0")); + } + + let rs = range.map(|value| match value { + Range::Int { first, last } => HTTPRangeSpec { + is_suffix_length: false, + start: first as i64, + end: last.map_or(-1, |last| last as i64), + }, + Range::Suffix { length } => HTTPRangeSpec { + is_suffix_length: true, + start: length as i64, + end: -1, + }, + }); + + if rs.is_some() && part_number.is_some() { + return Err(s3_error!(InvalidArgument, "range and part_number invalid")); + } + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), part_number, &req.headers) + .await + .map_err(ApiError::from)?; + + ( + GetObjectRequestContext { + bucket, + key, + part_number, + rs, + opts, + headers: req.headers.clone(), + sse_customer_key: req.input.sse_customer_key.clone(), + sse_customer_key_md5: req.input.sse_customer_key_md5.clone(), + }, + version_id.unwrap_or_default(), + ) + }; let base_buffer_size = self.base_buffer_size(); let manager = get_concurrency_manager(); - let flow_runtime = GetObjectFlowRuntime { - manager, - bootstrap: &bootstrap, - base_buffer_size, - }; - let helper = new_operation_helper(&req, EventName::ObjectAccessedGet, S3Operation::GetObject, true); - let flow_result = get_object_flow::run_get_object_flow(request_context.clone(), flow_runtime).await; + let cors_bucket = request_context.bucket.clone(); + let cors_method = req.method.clone(); + let cors_headers = request_context.headers.clone(); + let helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject).suppress_event(); + let flow_result = + get_object_flow::run_get_object_flow(request_context, version_id_for_event, manager, &bootstrap, base_buffer_size) + .await; let GetObjectBootstrap { mut request_guard, @@ -1042,7 +1167,15 @@ impl DefaultObjectUsecase { } = bootstrap; let result = match flow_result { - Ok(flow_result) => complete_get_flow_result(helper, &request_context, flow_result).await, + Ok(flow_result) => { + let helper = helper + .object(flow_result.event_info) + .version_id(flow_result.version_id_for_event); + let response = wrap_response_with_cors(&cors_bucket, &cors_method, &cors_headers, flow_result.output).await; + let result = Ok(response); + let _ = helper.complete(&result); + result + } Err(err) => Err(err), }; @@ -1766,7 +1899,7 @@ impl DefaultObjectUsecase { .await .map_err(ApiError::from)?; - maybe_enqueue_transition_immediate(&oi, LcEventSrc::S3CopyObject).await; + enqueue_transition_immediate(&oi, LcEventSrc::S3CopyObject).await; // Update quota tracking after successful copy if has_bucket_metadata { @@ -2975,8 +3108,19 @@ impl DefaultObjectUsecase { #[instrument(level = "debug", skip(self, req))] pub async fn execute_put_object_extract(&self, req: S3Request) -> S3Result> { - let request_context = prepare_put_object_request_context(&req); - let helper = new_operation_helper(&req, EventName::ObjectCreatedPut, S3Operation::PutObject, true); + let request_context = PutObjectRequestContext { + headers: req.headers.clone(), + trailing_headers: req.trailing_headers.clone(), + uri_query: req.uri.query().map(str::to_string), + is_post_object: req.extensions.get::().is_some(), + method: req.method.clone(), + uri: req.uri.clone(), + extensions: req.extensions.clone(), + credentials: req.credentials.clone(), + region: req.region.clone(), + service: req.service.clone(), + }; + let helper = OperationHelper::new(&req, EventName::ObjectCreatedPut, S3Operation::PutObject).suppress_event(); if is_sse_kms_requested(&req.input, &request_context.headers) { return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for extract uploads")); } @@ -2990,7 +3134,9 @@ impl DefaultObjectUsecase { .unwrap_or_else(default_notify_interface); let input = req.input; let output = DefaultObjectUsecase::run_put_object_extract_flow(input, request_context, notify, resolved_size).await?; - complete_put_response(helper, output) + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + result } } @@ -3022,37 +3168,6 @@ mod tests { } } - #[test] - fn put_object_execution_context_defaults_to_put() { - let input = PutObjectInput::builder() - .bucket("test-bucket".to_string()) - .key("test-key".to_string()) - .build() - .unwrap(); - let req = build_request(input, Method::PUT); - - let (event_name, quota_operation, method_name) = put_object_execution_context(&req); - assert_eq!(event_name, EventName::ObjectCreatedPut); - assert!(matches!(quota_operation, QuotaOperation::PutObject)); - assert_eq!(method_name, "PUT"); - } - - #[test] - fn put_object_execution_context_uses_post_marker() { - let input = PutObjectInput::builder() - .bucket("test-bucket".to_string()) - .key("test-key".to_string()) - .build() - .unwrap(); - let mut req = build_request(input, Method::POST); - req.extensions.insert(PostObjectRequestMarker); - - let (event_name, quota_operation, method_name) = put_object_execution_context(&req); - assert_eq!(event_name, EventName::ObjectCreatedPost); - assert!(matches!(quota_operation, QuotaOperation::PostObject)); - assert_eq!(method_name, "POST"); - } - #[tokio::test] async fn execute_put_object_rejects_invalid_storage_class() { let input = PutObjectInput::builder() diff --git a/rustfs/src/app/object_usecase/app_adapters.rs b/rustfs/src/app/object_usecase/app_adapters.rs deleted file mode 100644 index 268c10eb5..000000000 --- a/rustfs/src/app/object_usecase/app_adapters.rs +++ /dev/null @@ -1,461 +0,0 @@ -// 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 super::get_object_flow::GetObjectBootstrap; -use super::*; -use crate::app::context::NotifyInterface; -use crate::storage::concurrency::{self, ConcurrencyManager, get_buffer_size_opt_in}; -use hashbrown::HashMap; -use rustfs_object_io::get::{ - GetObjectBodyPlan as ObjectIoGetObjectBodyPlan, GetObjectBodyPlanningInputs as ObjectIoGetObjectBodyPlanningInputs, - GetObjectDataPlaneMetricContract as ObjectIoGetObjectDataPlaneMetricContract, GetObjectFlowResult, - MaterializeGetObjectBodyError as ObjectIoMaterializeGetObjectBodyError, - materialize_get_object_body as object_io_materialize_get_object_body, plan_get_object_body as object_io_plan_get_object_body, - plan_get_object_strategy_layout as object_io_plan_get_object_strategy_layout, -}; - -pub(super) async fn prepare_get_object_request_context(req: &S3Request) -> S3Result { - let GetObjectInput { - bucket, - key, - version_id, - part_number, - range, - .. - } = req.input.clone(); - - validate_object_key(&key, "GET")?; - - let part_number = part_number.map(|v| v as usize); - - if let Some(part_num) = part_number - && part_num == 0 - { - return Err(s3_error!(InvalidArgument, "Invalid part number: part number must be greater than 0")); - } - - let rs = range.map(|v| match v { - Range::Int { first, last } => HTTPRangeSpec { - is_suffix_length: false, - start: first as i64, - end: if let Some(last) = last { last as i64 } else { -1 }, - }, - Range::Suffix { length } => HTTPRangeSpec { - is_suffix_length: true, - start: length as i64, - end: -1, - }, - }); - - if rs.is_some() && part_number.is_some() { - return Err(s3_error!(InvalidArgument, "range and part_number invalid")); - } - - let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), part_number, &req.headers) - .await - .map_err(ApiError::from)?; - - Ok(GetObjectRequestContext { - version_id_for_event: version_id.unwrap_or_default(), - bucket, - key, - part_number, - rs, - opts, - headers: req.headers.clone(), - method: req.method.clone(), - sse_customer_key: req.input.sse_customer_key.clone(), - sse_customer_key_md5: req.input.sse_customer_key_md5.clone(), - }) -} - -pub(super) fn init_get_object_bootstrap(bucket: &str, key: &str, request_id: &str) -> S3Result { - let timeout_config = TimeoutConfig::from_env(); - let wrapper = RequestTimeoutWrapper::with_request_id(timeout_config.clone(), request_id.to_string()); - let request_start = std::time::Instant::now(); - let request_guard = ConcurrencyManager::track_request(); - let concurrent_requests = GetObjectGuard::concurrent_requests(); - - let deadlock_detector = deadlock_detector::get_deadlock_detector(); - deadlock_detector.register_request(request_id, format!("GetObject {bucket}/{key}")); - let deadlock_request_guard = DeadlockRequestGuard::new(deadlock_detector, request_id.to_string()); - - if wrapper.is_timeout() { - warn!( - bucket = %bucket, - key = %key, - timeout_secs = timeout_config.get_object_timeout.as_secs(), - elapsed_ms = wrapper.elapsed().as_millis(), - "GetObject request timed out before processing" - ); - return Err(s3_error!(InternalError, "Request timeout before processing")); - } - - rustfs_io_metrics::record_get_object_request_start(concurrent_requests); - - debug!( - "GetObject request started with {} concurrent requests, timeout={:?}", - concurrent_requests, timeout_config.get_object_timeout - ); - - Ok(GetObjectBootstrap { - timeout_config, - wrapper, - request_start, - request_guard, - _deadlock_request_guard: deadlock_request_guard, - }) -} - -pub(super) async fn build_get_object_body_adapter( - final_stream: R, - bucket: &str, - key: &str, - response_content_length: i64, - optimal_buffer_size: usize, - planning_inputs: ObjectIoGetObjectBodyPlanningInputs, -) -> S3Result> -where - R: AsyncRead + Send + Sync + Unpin + 'static, -{ - let body_plan = object_io_plan_get_object_body(planning_inputs, rustfs_config::DEFAULT_OBJECT_SEEK_SUPPORT_THRESHOLD); - - match body_plan { - ObjectIoGetObjectBodyPlan::BufferSeekable => { - debug!( - bucket = %bucket, - key = %key, - size = response_content_length, - "reading object into memory for seek support" - ); - } - ObjectIoGetObjectBodyPlan::Stream if planning_inputs.encryption_applied => { - info!( - "Encrypted object: Using unlimited stream for decryption with buffer size {}", - optimal_buffer_size - ); - } - _ => {} - } - - let materialized = - object_io_materialize_get_object_body(final_stream, body_plan, response_content_length, optimal_buffer_size) - .await - .map_err(|err| match err { - ObjectIoMaterializeGetObjectBodyError::EncryptedRead(err) => { - error!("Failed to read decrypted object into memory: {}", err); - ApiError::from(StorageError::other(format!("Failed to read decrypted object: {err}"))) - } - })?; - - Ok(materialized.body) -} - -pub(super) struct GetObjectCompletionInputs<'a> { - pub(super) bucket: &'a str, - pub(super) key: &'a str, - pub(super) wrapper: &'a RequestTimeoutWrapper, - pub(super) timeout_config: &'a TimeoutConfig, - pub(super) total_duration: Duration, - pub(super) response_content_length: i64, - pub(super) optimal_buffer_size: usize, - pub(super) metric_contract: ObjectIoGetObjectDataPlaneMetricContract, -} - -pub(super) struct GetObjectStrategyRuntimeInputs<'a> { - pub(super) base_buffer_size: usize, - pub(super) manager: &'a ConcurrencyManager, - pub(super) bucket: &'a str, - pub(super) key: &'a str, - pub(super) info: &'a ObjectInfo, - pub(super) rs: Option<&'a HTTPRangeSpec>, - pub(super) response_content_length: i64, - pub(super) permit_wait_duration: Duration, - pub(super) queue_utilization: f64, - pub(super) queue_status: &'a concurrency::IoQueueStatus, -} - -pub(super) fn finalize_get_object_completion(inputs: GetObjectCompletionInputs<'_>) { - let GetObjectCompletionInputs { - bucket, - key, - wrapper, - timeout_config, - total_duration, - response_content_length, - optimal_buffer_size, - metric_contract, - } = inputs; - - rustfs_io_metrics::record_get_object_completion(total_duration.as_secs_f64(), response_content_length, optimal_buffer_size); - - rustfs_io_metrics::record_get_object(total_duration.as_millis() as f64, response_content_length); - rustfs_io_metrics::record_io_copy_mode("get", metric_contract.copy_mode, response_content_length.max(0) as usize); - - if wrapper.is_timeout() { - warn!( - bucket = %bucket, - key = %key, - elapsed = ?wrapper.elapsed(), - timeout = ?timeout_config.get_object_timeout, - "GetObject request exceeded timeout" - ); - rustfs_io_metrics::record_get_object_timeout(None, Some(wrapper.elapsed().as_secs_f64())); - } - - debug!( - bucket = %bucket, - key = %key, - size = response_content_length, - duration = ?total_duration, - buffer = optimal_buffer_size, - "GetObject completed" - ); -} - -pub(super) fn finalize_get_object_strategy_runtime(inputs: GetObjectStrategyRuntimeInputs<'_>) -> usize { - let GetObjectStrategyRuntimeInputs { - base_buffer_size, - manager, - bucket, - key, - info, - rs, - response_content_length, - permit_wait_duration, - queue_utilization, - queue_status, - } = inputs; - - let strategy_layout = object_io_plan_get_object_strategy_layout( - rs, - response_content_length, - 0, - get_buffer_size_opt_in(response_content_length), - ); - - if let Some(range_spec) = rs - && range_spec.start >= 0 - { - manager.record_access(range_spec.start as u64, response_content_length as u64); - } - - if response_content_length > 0 { - manager.record_transfer(response_content_length as u64, permit_wait_duration); - } - - let io_strategy = manager.calculate_io_strategy_with_context( - info.size, - base_buffer_size, - permit_wait_duration, - strategy_layout.is_sequential_hint, - ); - - debug!( - wait_ms = permit_wait_duration.as_millis() as u64, - load_level = ?io_strategy.load_level, - buffer_size = io_strategy.buffer_size, - buffer_multiplier = io_strategy.buffer_multiplier, - readahead = io_strategy.enable_readahead, - storage_media = ?io_strategy.storage_media, - access_pattern = ?io_strategy.access_pattern, - bandwidth_tier = ?io_strategy.bandwidth_tier, - concurrent_requests = io_strategy.concurrent_requests, - file_size = info.size, - is_sequential = strategy_layout.is_sequential_hint, - "Enhanced multi-factor I/O strategy calculated" - ); - - let io_priority = manager.get_io_priority(response_content_length); - - if manager.is_priority_scheduling_enabled() { - debug!( - bucket = %bucket, - key = %key, - priority = %io_priority, - request_size = response_content_length, - "I/O priority assigned (based on actual request size)" - ); - - rustfs_io_metrics::record_io_priority_assignment(io_priority.as_str()); - } - - rustfs_io_metrics::record_get_object_io_state( - permit_wait_duration.as_secs_f64(), - queue_utilization, - queue_status.permits_in_use, - queue_status.total_permits.saturating_sub(queue_status.permits_in_use), - io_strategy.load_level.as_str(), - io_strategy.buffer_multiplier, - ); - - let strategy_layout = object_io_plan_get_object_strategy_layout( - rs, - response_content_length, - io_strategy.buffer_size, - get_buffer_size_opt_in(response_content_length), - ); - - debug!( - actual_request_size = response_content_length, - priority = %io_priority.as_str(), - "I/O priority finalized with actual request size" - ); - - debug!( - "GetObject buffer sizing: file_size={}, base={}, optimal={}, concurrent_requests={}, io_strategy={:?}", - response_content_length, - get_buffer_size_opt_in(response_content_length), - strategy_layout.optimal_buffer_size, - io_strategy.concurrent_requests, - io_strategy.load_level - ); - - strategy_layout.optimal_buffer_size -} - -pub(super) fn prepare_put_object_request_context(req: &S3Request) -> PutObjectRequestContext { - PutObjectRequestContext { - headers: req.headers.clone(), - trailing_headers: req.trailing_headers.clone(), - uri_query: req.uri.query().map(str::to_string), - is_post_object: req.extensions.get::().is_some(), - method: req.method.clone(), - uri: req.uri.clone(), - extensions: req.extensions.clone(), - credentials: req.credentials.clone(), - region: req.region.clone(), - service: req.service.clone(), - } -} - -pub(super) fn put_object_execution_context(req: &S3Request) -> (EventName, QuotaOperation, &'static str) { - if req.extensions.get::().is_some() { - (EventName::ObjectCreatedPost, QuotaOperation::PostObject, "POST") - } else { - (EventName::ObjectCreatedPut, QuotaOperation::PutObject, "PUT") - } -} - -pub(super) fn new_operation_helper( - req: &S3Request, - event_name: EventName, - operation: S3Operation, - suppress_event: bool, -) -> OperationHelper { - let helper = OperationHelper::new(req, event_name, operation); - if suppress_event { helper.suppress_event() } else { helper } -} - -pub(super) fn bind_helper_object( - helper: OperationHelper, - object_info: ObjectInfo, - version_id: Option, -) -> OperationHelper { - let helper = helper.object(object_info); - if let Some(version_id) = version_id { - helper.version_id(version_id) - } else { - helper - } -} - -pub(super) async fn complete_get_flow_result( - helper: OperationHelper, - request_context: &GetObjectRequestContext, - flow_result: GetObjectFlowResult, -) -> S3Result> { - let helper = helper - .object(flow_result.event_info) - .version_id(flow_result.version_id_for_event); - let response = wrap_response_with_cors( - &request_context.bucket, - &request_context.method, - &request_context.headers, - flow_result.output, - ) - .await; - let result = Ok(response); - let _ = helper.complete(&result); - result -} - -pub(super) fn complete_put_response(helper: OperationHelper, output: PutObjectOutput) -> S3Result> { - let result = Ok(S3Response::new(output)); - let _ = helper.complete(&result); - result -} - -#[allow(clippy::too_many_arguments)] -pub(super) fn spawn_put_extract_notification( - notify: Arc, - request_context: Option, - bucket: String, - req_params: HashMap, - version_id: String, - host: String, - port: u16, - user_agent: String, - obj_info: ObjectInfo, - output: PutObjectOutput, -) { - let event_args = rustfs_notify::EventArgs { - event_name: EventName::ObjectCreatedPut, - bucket_name: bucket, - object: obj_info, - req_params, - resp_elements: extract_resp_elements(&S3Response::new(output)), - version_id, - host, - port, - user_agent, - }; - - helper::spawn_background_with_context(request_context, async move { - notify.notify(event_args).await; - }); -} - -pub(super) async fn get_validated_store_adapter(bucket: &str) -> S3Result> { - get_validated_store(bucket).await -} - -pub(super) async fn bucket_prefix_versioning_enabled(bucket: &str, key: &str) -> bool { - BucketVersioningSys::prefix_enabled(bucket, key).await -} - -pub(super) async fn authorize_extract_put_target( - request_context: &PutObjectRequestContext, - bucket: &str, - object: &str, -) -> S3Result<()> { - let mut auth_req = S3Request { - input: PutObjectInput::default(), - method: request_context.method.clone(), - uri: request_context.uri.clone(), - headers: request_context.headers.clone(), - extensions: request_context.extensions.clone(), - credentials: request_context.credentials.clone(), - region: request_context.region.clone(), - service: request_context.service.clone(), - trailing_headers: request_context.trailing_headers.clone(), - }; - { - let req_info = req_info_mut(&mut auth_req)?; - req_info.bucket = Some(bucket.to_string()); - req_info.object = Some(object.to_string()); - req_info.version_id = None; - } - authorize_request(&mut auth_req, Action::S3Action(S3Action::PutObjectAction)).await -} diff --git a/rustfs/src/app/object_usecase/get_object_flow.rs b/rustfs/src/app/object_usecase/get_object_flow.rs index 627425895..98b63d871 100644 --- a/rustfs/src/app/object_usecase/get_object_flow.rs +++ b/rustfs/src/app/object_usecase/get_object_flow.rs @@ -13,29 +13,32 @@ // limitations under the License. use super::DeadlockRequestGuard; -use super::app_adapters::{ - GetObjectCompletionInputs, GetObjectStrategyRuntimeInputs, bucket_prefix_versioning_enabled, build_get_object_body_adapter, - finalize_get_object_completion, finalize_get_object_strategy_runtime, -}; -use super::get_object_zero_copy::{GetObjectPreparedRead, prepare_get_object_read_execution}; +use super::get_object_zero_copy::{GetObjectIoPlanning, GetObjectPreparedRead, prepare_get_object_read_execution}; use super::types::GetObjectRequestContext; use crate::error::ApiError; -use crate::storage::concurrency::{self, ConcurrencyManager, GetObjectGuard}; +use crate::storage::concurrency::{ConcurrencyManager, GetObjectGuard, get_buffer_size_opt_in}; use crate::storage::options::filter_object_metadata; use crate::storage::timeout_wrapper::{RequestTimeoutWrapper, TimeoutConfig}; +use rustfs_ecstore::bucket::versioning_sys::BucketVersioningSys; +use rustfs_ecstore::error::StorageError; use rustfs_ecstore::store_api::{HTTPRangeSpec, ObjectInfo}; use rustfs_object_io::get::{ - GetObjectBodyPlanningInputs as ObjectIoGetObjectBodyPlanningInputs, GetObjectBodySource, - GetObjectDataPlaneMetricContract as ObjectIoGetObjectDataPlaneMetricContract, GetObjectFlowResult, GetObjectOutputContext, - GetObjectReadSetup, build_chunk_blob as object_io_build_chunk_blob, + GetObjectBodyPlan as ObjectIoGetObjectBodyPlan, GetObjectBodyPlanningInputs as ObjectIoGetObjectBodyPlanningInputs, + GetObjectBodySource, GetObjectDataPlaneMetricContract as ObjectIoGetObjectDataPlaneMetricContract, GetObjectFlowResult, + GetObjectOutputContext, GetObjectReadSetup, MaterializeGetObjectBodyError as ObjectIoMaterializeGetObjectBodyError, + build_chunk_blob as object_io_build_chunk_blob, build_cors_wrapped_get_object_flow_result as object_io_build_cors_wrapped_get_object_flow_result, build_get_object_checksums as object_io_build_get_object_checksums, build_get_object_output_context as object_io_build_get_object_output_context, chunk_body_data_plane_labels as object_io_chunk_body_data_plane_labels, + materialize_get_object_body as object_io_materialize_get_object_body, plan_get_object_body as object_io_plan_get_object_body, + plan_get_object_strategy_layout as object_io_plan_get_object_strategy_layout, }; use s3s::S3Result; -use s3s::dto::{ContentType, SSECustomerAlgorithm, SSECustomerKeyMD5, SSEKMSKeyId, ServerSideEncryption, Timestamp}; +use s3s::dto::StreamingBlob; use std::time::Duration; +use tokio::io::AsyncRead; +use tracing::{debug, error, info, warn}; pub(super) struct GetObjectBootstrap { pub(super) timeout_config: TimeoutConfig, @@ -45,51 +48,228 @@ pub(super) struct GetObjectBootstrap { pub(super) _deadlock_request_guard: DeadlockRequestGuard, } -#[derive(Clone, Copy)] -pub(super) struct GetObjectFlowRuntime<'a> { - pub(super) manager: &'a ConcurrencyManager, - pub(super) bootstrap: &'a GetObjectBootstrap, - pub(super) base_buffer_size: usize, +async fn build_get_object_body_adapter( + final_stream: R, + bucket: &str, + key: &str, + response_content_length: i64, + optimal_buffer_size: usize, + planning_inputs: ObjectIoGetObjectBodyPlanningInputs, +) -> S3Result> +where + R: AsyncRead + Send + Sync + Unpin + 'static, +{ + let body_plan = object_io_plan_get_object_body(planning_inputs, rustfs_config::DEFAULT_OBJECT_SEEK_SUPPORT_THRESHOLD); + + match body_plan { + ObjectIoGetObjectBodyPlan::BufferSeekable => { + debug!( + bucket = %bucket, + key = %key, + size = response_content_length, + "reading object into memory for seek support" + ); + } + ObjectIoGetObjectBodyPlan::Stream if planning_inputs.encryption_applied => { + info!( + "Encrypted object: Using unlimited stream for decryption with buffer size {}", + optimal_buffer_size + ); + } + _ => {} + } + + let materialized = + object_io_materialize_get_object_body(final_stream, body_plan, response_content_length, optimal_buffer_size) + .await + .map_err(|err| match err { + ObjectIoMaterializeGetObjectBodyError::EncryptedRead(err) => { + error!("Failed to read decrypted object into memory: {}", err); + ApiError::from(StorageError::other(format!("Failed to read decrypted object: {err}"))) + } + })?; + + Ok(materialized.body) +} + +fn finalize_get_object_completion( + request_context: &GetObjectRequestContext, + wrapper: &RequestTimeoutWrapper, + timeout_config: &TimeoutConfig, + total_duration: Duration, + response_content_length: i64, + optimal_buffer_size: usize, + metric_contract: ObjectIoGetObjectDataPlaneMetricContract, +) { + rustfs_io_metrics::record_get_object_completion(total_duration.as_secs_f64(), response_content_length, optimal_buffer_size); + + rustfs_io_metrics::record_get_object(total_duration.as_millis() as f64, response_content_length); + rustfs_io_metrics::record_io_copy_mode("get", metric_contract.copy_mode, response_content_length.max(0) as usize); + + if wrapper.is_timeout() { + warn!( + bucket = %request_context.bucket, + key = %request_context.key, + elapsed = ?wrapper.elapsed(), + timeout = ?timeout_config.get_object_timeout, + "GetObject request exceeded timeout" + ); + rustfs_io_metrics::record_get_object_timeout(None, Some(wrapper.elapsed().as_secs_f64())); + } + + debug!( + bucket = %request_context.bucket, + key = %request_context.key, + size = response_content_length, + duration = ?total_duration, + buffer = optimal_buffer_size, + "GetObject completed" + ); +} + +fn get_object_strategy_range<'a>( + request_context: &'a GetObjectRequestContext, + resolved_range: Option<&'a HTTPRangeSpec>, +) -> Option<&'a HTTPRangeSpec> { + resolved_range.or(request_context.rs.as_ref()) +} + +fn finalize_get_object_strategy_runtime( + request_context: &GetObjectRequestContext, + resolved_range: Option<&HTTPRangeSpec>, + manager: &ConcurrencyManager, + base_buffer_size: usize, + info: &ObjectInfo, + response_content_length: i64, + io_planning: &GetObjectIoPlanning<'_>, +) -> usize { + let strategy_range = get_object_strategy_range(request_context, resolved_range); + let strategy_layout = object_io_plan_get_object_strategy_layout( + strategy_range, + response_content_length, + 0, + get_buffer_size_opt_in(response_content_length), + ); + + if let Some(range_spec) = strategy_range + && range_spec.start >= 0 + { + manager.record_access(range_spec.start as u64, response_content_length as u64); + } + + if response_content_length > 0 { + manager.record_transfer(response_content_length as u64, io_planning.permit_wait_duration); + } + + let io_strategy = manager.calculate_io_strategy_with_context( + info.size, + base_buffer_size, + io_planning.permit_wait_duration, + strategy_layout.is_sequential_hint, + ); + + debug!( + wait_ms = io_planning.permit_wait_duration.as_millis() as u64, + load_level = ?io_strategy.load_level, + buffer_size = io_strategy.buffer_size, + buffer_multiplier = io_strategy.buffer_multiplier, + readahead = io_strategy.enable_readahead, + storage_media = ?io_strategy.storage_media, + access_pattern = ?io_strategy.access_pattern, + bandwidth_tier = ?io_strategy.bandwidth_tier, + concurrent_requests = io_strategy.concurrent_requests, + file_size = info.size, + is_sequential = strategy_layout.is_sequential_hint, + "Enhanced multi-factor I/O strategy calculated" + ); + + let io_priority = manager.get_io_priority(response_content_length); + + if manager.is_priority_scheduling_enabled() { + debug!( + bucket = %request_context.bucket, + key = %request_context.key, + priority = %io_priority, + request_size = response_content_length, + "I/O priority assigned (based on actual request size)" + ); + + rustfs_io_metrics::record_io_priority_assignment(io_priority.as_str()); + } + + rustfs_io_metrics::record_get_object_io_state( + io_planning.permit_wait_duration.as_secs_f64(), + io_planning.queue_utilization, + io_planning.queue_status.permits_in_use, + io_planning + .queue_status + .total_permits + .saturating_sub(io_planning.queue_status.permits_in_use), + io_strategy.load_level.as_str(), + io_strategy.buffer_multiplier, + ); + + let strategy_layout = object_io_plan_get_object_strategy_layout( + strategy_range, + response_content_length, + io_strategy.buffer_size, + get_buffer_size_opt_in(response_content_length), + ); + + debug!( + actual_request_size = response_content_length, + priority = %io_priority.as_str(), + "I/O priority finalized with actual request size" + ); + + debug!( + "GetObject buffer sizing: file_size={}, base={}, optimal={}, concurrent_requests={}, io_strategy={:?}", + response_content_length, + get_buffer_size_opt_in(response_content_length), + strategy_layout.optimal_buffer_size, + io_strategy.concurrent_requests, + io_strategy.load_level + ); + + strategy_layout.optimal_buffer_size } -#[allow(clippy::too_many_arguments)] pub(super) async fn build_get_object_output_context( request_context: &GetObjectRequestContext, manager: &ConcurrencyManager, - bucket: &str, - key: &str, - info: ObjectInfo, - event_info: ObjectInfo, - body_source: GetObjectBodySource, - rs: Option, - content_type: Option, - last_modified: Option, - response_content_length: i64, - content_range: Option, - server_side_encryption: Option, - sse_customer_algorithm: Option, - sse_customer_key_md5: Option, - ssekms_key_id: Option, - encryption_applied: bool, - permit_wait_duration: Duration, - queue_utilization: f64, - queue_status: &concurrency::IoQueueStatus, + read_setup: GetObjectReadSetup, + io_planning: &GetObjectIoPlanning<'_>, base_buffer_size: usize, - part_number: Option, versioned: bool, ) -> S3Result<(GetObjectOutputContext, ObjectIoGetObjectDataPlaneMetricContract)> { - let optimal_buffer_size = finalize_get_object_strategy_runtime(GetObjectStrategyRuntimeInputs { - base_buffer_size, - manager, - bucket, - key, - info: &info, - rs: rs.as_ref(), + let bucket = &request_context.bucket; + let key = &request_context.key; + let part_number = request_context.part_number; + let GetObjectReadSetup { + info, + event_info, + body_source, + rs, + content_type, + last_modified, response_content_length, - permit_wait_duration, - queue_utilization, - queue_status, - }); + content_range, + server_side_encryption, + sse_customer_algorithm, + sse_customer_key_md5, + ssekms_key_id, + encryption_applied, + } = read_setup; + + let optimal_buffer_size = finalize_get_object_strategy_runtime( + request_context, + rs.as_ref(), + manager, + base_buffer_size, + &info, + response_content_length, + io_planning, + ); let (body, metric_contract) = match body_source { GetObjectBodySource::Reader(final_stream) => { @@ -156,97 +336,85 @@ pub(super) async fn build_get_object_output_context( pub(super) async fn run_get_object_flow( request_context: GetObjectRequestContext, - runtime: GetObjectFlowRuntime<'_>, + version_id_for_event: String, + manager: &ConcurrencyManager, + bootstrap: &GetObjectBootstrap, + base_buffer_size: usize, ) -> S3Result { - let GetObjectFlowRuntime { - manager, - bootstrap, - base_buffer_size, - } = runtime; let timeout_config = &bootstrap.timeout_config; let wrapper = &bootstrap.wrapper; let request_start = bootstrap.request_start; - let bucket = request_context.bucket.clone(); - let key = request_context.key.clone(); - let version_id_for_event = request_context.version_id_for_event.clone(); - let part_number = request_context.part_number; - let rs = request_context.rs.clone(); - let opts = request_context.opts.clone(); - let prepared_read = prepare_get_object_read_execution( - &request_context, - manager, - wrapper, - timeout_config, - &bucket, - &key, - rs, - &opts, - part_number, - ) - .await?; + let prepared_read = prepare_get_object_read_execution(&request_context, manager, wrapper, timeout_config).await?; let GetObjectPreparedRead { io_planning, read_setup } = prepared_read; - let permit_wait_duration = io_planning.permit_wait_duration; - let queue_status = io_planning.queue_status; - let queue_utilization = io_planning.queue_utilization; - let GetObjectReadSetup { - info, - event_info, - body_source, - rs, - content_type, - last_modified, - response_content_length, - content_range, - server_side_encryption, - sse_customer_algorithm, - sse_customer_key_md5, - ssekms_key_id, - encryption_applied, - } = read_setup; - - let versioned = bucket_prefix_versioning_enabled(&bucket, &key).await; - let (output_context, metric_contract) = build_get_object_output_context( - &request_context, - manager, - &bucket, - &key, - info, - event_info, - body_source, - rs, - content_type, - last_modified, - response_content_length, - content_range, - server_side_encryption, - sse_customer_algorithm, - sse_customer_key_md5, - ssekms_key_id, - encryption_applied, - permit_wait_duration, - queue_utilization, - &queue_status, - base_buffer_size, - part_number, - versioned, - ) - .await?; + let versioned = BucketVersioningSys::prefix_enabled(&request_context.bucket, &request_context.key).await; + let (output_context, metric_contract) = + build_get_object_output_context(&request_context, manager, read_setup, &io_planning, base_buffer_size, versioned).await?; let response_content_length = output_context.response_content_length; let optimal_buffer_size = output_context.optimal_buffer_size; let total_duration = request_start.elapsed(); - finalize_get_object_completion(GetObjectCompletionInputs { - bucket: &bucket, - key: &key, + finalize_get_object_completion( + &request_context, wrapper, timeout_config, total_duration, response_content_length, optimal_buffer_size, metric_contract, - }); + ); Ok(object_io_build_cors_wrapped_get_object_flow_result(output_context, version_id_for_event)) } + +#[cfg(test)] +mod tests { + use super::get_object_strategy_range; + use super::*; + use http::HeaderMap; + use rustfs_ecstore::store_api::ObjectOptions; + + fn sample_range(start: i64, end: i64) -> HTTPRangeSpec { + HTTPRangeSpec { + is_suffix_length: false, + start, + end, + } + } + + fn sample_request_context() -> GetObjectRequestContext { + GetObjectRequestContext { + bucket: "bucket".to_string(), + key: "key".to_string(), + part_number: None, + rs: None, + opts: ObjectOptions::default(), + headers: HeaderMap::new(), + sse_customer_key: None, + sse_customer_key_md5: None, + } + } + + #[test] + fn strategy_range_prefers_resolved_range_for_part_reads() { + let request_context = sample_request_context(); + let resolved_range = sample_range(1024, 2047); + + let strategy_range = get_object_strategy_range(&request_context, Some(&resolved_range)).unwrap(); + + assert_eq!(strategy_range.start, 1024); + assert_eq!(strategy_range.end, 2047); + } + + #[test] + fn strategy_range_falls_back_to_raw_request_range() { + let mut request_context = sample_request_context(); + request_context.rs = Some(sample_range(0, 511)); + + let strategy_range = get_object_strategy_range(&request_context, None).unwrap(); + + assert_eq!(strategy_range.start, 0); + assert_eq!(strategy_range.end, 511); + } +} diff --git a/rustfs/src/app/object_usecase/get_object_zero_copy.rs b/rustfs/src/app/object_usecase/get_object_zero_copy.rs index 0df64b82c..cc05b5b48 100644 --- a/rustfs/src/app/object_usecase/get_object_zero_copy.rs +++ b/rustfs/src/app/object_usecase/get_object_zero_copy.rs @@ -12,17 +12,17 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::app_adapters::get_validated_store_adapter; use super::types::GetObjectRequestContext; use crate::error::ApiError; use crate::storage::concurrency::{self, ConcurrencyManager}; use crate::storage::timeout_wrapper::{RequestTimeoutWrapper, TimeoutConfig}; use crate::storage::{ - DecryptionRequest, check_preconditions, sse_decryption, validate_sse_headers_for_read, validate_ssec_for_read, + DecryptionRequest, check_preconditions, get_validated_store, sse_decryption, validate_sse_headers_for_read, + validate_ssec_for_read, }; use http::HeaderMap; use rustfs_concurrency::GetObjectQueueSnapshot; -use rustfs_ecstore::store_api::{HTTPRangeSpec, ObjectIO, ObjectOperations, ObjectOptions}; +use rustfs_ecstore::store_api::{ObjectIO, ObjectOperations}; use rustfs_object_io::get::{ ChunkReadDecision, ChunkReadPlanError, GetObjectEncryptionState as ObjectIoGetObjectEncryptionState, GetObjectReadSetup, build_reader_read_setup as object_io_build_reader_read_setup, @@ -115,21 +115,20 @@ pub(super) async fn acquire_get_object_io_planning<'a>( }) } -#[allow(clippy::too_many_arguments)] pub(super) async fn prepare_get_object_read( request_context: &GetObjectRequestContext, store: &rustfs_ecstore::store::ECStore, manager: &ConcurrencyManager, - bucket: &str, - key: &str, - rs: Option, - h: HeaderMap, - opts: &ObjectOptions, - part_number: Option, read_start: std::time::Instant, ) -> S3Result { let reader = store - .get_object_reader(bucket, key, rs.clone(), h, opts) + .get_object_reader( + &request_context.bucket, + &request_context.key, + request_context.rs.clone(), + HeaderMap::new(), + &request_context.opts, + ) .await .map_err(ApiError::from)?; @@ -159,7 +158,8 @@ pub(super) async fn prepare_get_object_read( request_context.sse_customer_key.as_ref(), request_context.sse_customer_key_md5.as_ref(), )?; - let read_plan = object_io_plan_legacy_read(&info, rs, part_number).map_err(ApiError::from)?; + let read_plan = + object_io_plan_legacy_read(&info, request_context.rs.clone(), request_context.part_number).map_err(ApiError::from)?; debug!( "GET object metadata check: parts={}, provided_sse_key={:?}", @@ -168,8 +168,8 @@ pub(super) async fn prepare_get_object_read( ); let decryption_request = DecryptionRequest { - bucket, - key, + bucket: &request_context.bucket, + key: &request_context.key, metadata: &info.user_defined, sse_customer_key: request_context.sse_customer_key.as_ref(), sse_customer_key_md5: request_context.sse_customer_key_md5.as_ref(), @@ -218,68 +218,44 @@ pub(super) async fn prepare_get_object_read( )) } -#[allow(clippy::too_many_arguments)] pub(super) async fn prepare_get_object_read_execution<'a>( request_context: &GetObjectRequestContext, manager: &'a ConcurrencyManager, wrapper: &RequestTimeoutWrapper, timeout_config: &TimeoutConfig, - bucket: &str, - key: &str, - rs: Option, - opts: &ObjectOptions, - part_number: Option, ) -> S3Result> { - let h = HeaderMap::new(); - let io_planning = acquire_get_object_io_planning(manager, wrapper, timeout_config, bucket, key).await?; - let store = get_validated_store_adapter(bucket).await?; + let io_planning = + acquire_get_object_io_planning(manager, wrapper, timeout_config, &request_context.bucket, &request_context.key).await?; + let store = get_validated_store(&request_context.bucket).await?; let read_start = std::time::Instant::now(); let read_setup = match object_io_get_object_chunk_fast_path_guard( request_context.sse_customer_key.is_some(), request_context.sse_customer_key_md5.is_some(), ) { - Ok(()) => match prepare_get_object_chunk_read( - request_context, - &store, - manager, - bucket, - key, - rs.clone(), - part_number, - opts, - read_start, - ) - .await? - { + Ok(()) => match prepare_get_object_chunk_read(request_context, &store, manager, read_start).await? { Some(read_setup) => read_setup, - None => { - prepare_get_object_read(request_context, &store, manager, bucket, key, rs, h, opts, part_number, read_start) - .await? - } + None => prepare_get_object_read(request_context, &store, manager, read_start).await?, }, Err(fallback) => { rustfs_io_metrics::record_io_fallback(fallback.stage, fallback.reason); - prepare_get_object_read(request_context, &store, manager, bucket, key, rs, h, opts, part_number, read_start).await? + prepare_get_object_read(request_context, &store, manager, read_start).await? } }; Ok(GetObjectPreparedRead { io_planning, read_setup }) } -#[allow(clippy::too_many_arguments)] pub(super) async fn prepare_get_object_chunk_read( request_context: &GetObjectRequestContext, store: &rustfs_ecstore::store::ECStore, manager: &ConcurrencyManager, - bucket: &str, - key: &str, - mut rs: Option, - part_number: Option, - opts: &ObjectOptions, read_start: std::time::Instant, ) -> S3Result> { - let info = store.get_object_info(bucket, key, opts).await.map_err(ApiError::from)?; + let info = store + .get_object_info(&request_context.bucket, &request_context.key, &request_context.opts) + .await + .map_err(ApiError::from)?; validate_sse_headers_for_read(&info.user_defined, &request_context.headers)?; validate_ssec_for_read( @@ -301,7 +277,12 @@ pub(super) async fn prepare_get_object_chunk_read( return Ok(None); } - let plan = match object_io_plan_chunk_read(&info, opts.version_id.is_none(), rs.clone(), part_number) { + let plan = match object_io_plan_chunk_read( + &info, + request_context.opts.version_id.is_none(), + request_context.rs.clone(), + request_context.part_number, + ) { Ok(ChunkReadDecision::Eligible(plan)) => plan, Ok(ChunkReadDecision::Fallback(fallback)) => { rustfs_io_metrics::record_io_fallback(fallback.stage, fallback.reason); @@ -311,14 +292,20 @@ pub(super) async fn prepare_get_object_chunk_read( Err(ChunkReadPlanError::MethodNotAllowed) => return Err(S3Error::new(S3ErrorCode::MethodNotAllowed)), Err(ChunkReadPlanError::Io(err)) => return Err(ApiError::from(err).into()), }; - rs = plan.rs.clone(); + let rs = plan.rs.clone(); let read_duration = read_start.elapsed(); manager.record_disk_operation(info.size as u64, read_duration, true).await; let event_info = info.clone(); let chunk_result = match store - .get_object_chunks(bucket, key, rs.clone(), HeaderMap::new(), opts) + .get_object_chunks( + &request_context.bucket, + &request_context.key, + rs.clone(), + HeaderMap::new(), + &request_context.opts, + ) .await .map_err(ApiError::from) { diff --git a/rustfs/src/app/object_usecase/put_object_extract.rs b/rustfs/src/app/object_usecase/put_object_extract.rs index 697f03880..a7806dc2d 100644 --- a/rustfs/src/app/object_usecase/put_object_extract.rs +++ b/rustfs/src/app/object_usecase/put_object_extract.rs @@ -64,29 +64,9 @@ impl DefaultObjectUsecase { let sse_customer_key_md5 = sse_customer_key_md5.or(h_md5); let original_sse = server_side_encryption.or(extract_server_side_encryption_from_headers(&request_context.headers)?); - let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok(); - let mut effective_sse = original_sse.or_else(|| { - bucket_sse_config.as_ref().and_then(|(config, _timestamp)| { - config.rules.first().and_then(|rule| { - rule.apply_server_side_encryption_by_default - .as_ref() - .map(|sse| match sse.sse_algorithm.as_str() { - "AES256" => ServerSideEncryption::from_static(ServerSideEncryption::AES256), - "aws:kms" => ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS), - _ => ServerSideEncryption::from_static(ServerSideEncryption::AES256), - }) - }) - }) - }); - let mut effective_kms_key_id = ssekms_key_id.or_else(|| { - bucket_sse_config.as_ref().and_then(|(config, _timestamp)| { - config.rules.first().and_then(|rule| { - rule.apply_server_side_encryption_by_default - .as_ref() - .and_then(|sse| sse.kms_master_key_id.clone()) - }) - }) - }); + let (default_sse, default_kms_key_id) = resolve_bucket_default_server_side_encryption(&bucket).await; + let mut effective_sse = original_sse.or(default_sse); + let mut effective_kms_key_id = ssekms_key_id.or(default_kms_key_id); if effective_sse .as_ref() .is_some_and(|sse| sse.as_str().eq_ignore_ascii_case(ServerSideEncryption::AWS_KMS)) @@ -196,7 +176,24 @@ impl DefaultObjectUsecase { let is_dir = f.header().entry_type().is_dir(); let fpath = normalize_extract_entry_key(&fpath.to_string_lossy(), extract_options.prefix.as_deref(), is_dir); - authorize_extract_put_target(&request_context, &bucket, &fpath).await?; + let mut auth_req = S3Request { + input: PutObjectInput::default(), + method: request_context.method.clone(), + uri: request_context.uri.clone(), + headers: request_context.headers.clone(), + extensions: request_context.extensions.clone(), + credentials: request_context.credentials.clone(), + region: request_context.region.clone(), + service: request_context.service.clone(), + trailing_headers: request_context.trailing_headers.clone(), + }; + { + let req_info = req_info_mut(&mut auth_req)?; + req_info.bucket = Some(bucket.clone()); + req_info.object = Some(fpath.clone()); + req_info.version_id = None; + } + authorize_request(&mut auth_req, Action::S3Action(S3Action::PutObjectAction)).await?; let mut size = f.header().size().unwrap_or_default() as i64; let archive_entry_mod_time = f @@ -319,18 +316,23 @@ impl DefaultObjectUsecase { ..Default::default() }; - spawn_put_extract_notification( - notify.clone(), - tracing_context.clone(), - bucket.clone(), - req_params.clone(), - version_id.clone(), - host.clone(), + let event_args = rustfs_notify::EventArgs { + event_name: EventName::ObjectCreatedPut, + bucket_name: bucket.clone(), + object: obj_info.clone(), + req_params: req_params.clone(), + resp_elements: extract_resp_elements(&S3Response::new(output)), + version_id: version_id.clone(), + host: host.clone(), port, - user_agent.clone(), - obj_info.clone(), - output, - ); + user_agent: user_agent.clone(), + }; + crate::storage::helper::spawn_background_with_context(tracing_context.clone(), { + let notify = notify.clone(); + async move { + notify.notify(event_args).await; + } + }); } let mut checksums = PutObjectChecksums { diff --git a/rustfs/src/app/object_usecase/put_object_flow.rs b/rustfs/src/app/object_usecase/put_object_flow.rs index 9ae0e5a8d..500c4ddc2 100644 --- a/rustfs/src/app/object_usecase/put_object_flow.rs +++ b/rustfs/src/app/object_usecase/put_object_flow.rs @@ -52,14 +52,6 @@ fn clamp_small_put_eager_max_bytes(inline_object_limit_bytes: Option) -> .min(DEFAULT_SMALL_PUT_EAGER_MAX_BYTES as usize) as i64 } -fn env_flag_enabled(name: &str) -> bool { - rustfs_utils::get_env_bool(name, false) -} - -fn env_non_negative_i64(name: &str) -> Option { - rustfs_utils::get_env_opt_i64(name).filter(|value| *value >= 0) -} - fn topology_aware_small_put_eager_max_bytes(store: &rustfs_ecstore::store::ECStore, versioned: bool) -> i64 { let Some(first_pool) = store.pools.first() else { return DEFAULT_SMALL_PUT_EAGER_MAX_BYTES; @@ -78,11 +70,12 @@ fn topology_aware_small_put_eager_max_bytes(store: &rustfs_ecstore::store::ECSto } fn resolved_small_put_eager_max_bytes(default_max_bytes: i64) -> i64 { - if env_flag_enabled(ENV_RUSTFS_PUT_FORCE_DISABLE_SMALL_EAGER) { + if rustfs_utils::get_env_bool(ENV_RUSTFS_PUT_FORCE_DISABLE_SMALL_EAGER, false) { return 0; } - env_non_negative_i64(ENV_RUSTFS_PUT_SMALL_EAGER_MAX_BYTES) + rustfs_utils::get_env_opt_i64(ENV_RUSTFS_PUT_SMALL_EAGER_MAX_BYTES) + .filter(|value| *value >= 0) .map(|value| value.min(DEFAULT_SMALL_PUT_EAGER_MAX_BYTES).min(default_max_bytes)) .unwrap_or(default_max_bytes) } @@ -91,37 +84,13 @@ fn should_use_small_put_eager_path(size: i64, eager_max_bytes: i64, compression_ size > 0 && size <= eager_max_bytes && !compression_enabled && !encryption_enabled } -fn request_uses_trailing_checksum(headers: &HeaderMap, trailing_headers: &Option) -> bool { - trailing_headers.is_some() - || headers.contains_key(AMZ_TRAILER) - || matches!( - rustfs_rio::get_content_checksum(headers), - Ok(Some(checksum)) if checksum.checksum_type.trailing() - ) -} - -fn put_path_label(small_eager: bool, reduced_copy: bool, compressed: bool) -> &'static str { - if small_eager { - "small_eager" - } else if compressed { - "compressed" - } else if reduced_copy { - "reduced_copy" - } else { - "legacy_plain" - } -} - -#[allow(clippy::too_many_arguments)] fn log_put_flow_phase( bucket: &str, key: &str, phase: &str, elapsed: Duration, object_size: i64, - small_eager: bool, - reduced_copy: bool, - compressed: bool, + put_path: &'static str, encrypted: bool, ) { let duration_ms = elapsed.as_millis() as u64; @@ -129,21 +98,20 @@ fn log_put_flow_phase( return; } - let put_path = put_path_label(small_eager, reduced_copy, compressed); if duration_ms >= SLOW_PUT_PHASE_ERROR_THRESHOLD_MS { error!( phase, - duration_ms, object_size, put_path, compressed, encrypted, bucket, key, "Small PUT phase is critically slow" + duration_ms, object_size, put_path, encrypted, bucket, key, "Small PUT phase is critically slow" ); } else if duration_ms >= SLOW_PUT_PHASE_WARN_THRESHOLD_MS { warn!( phase, - duration_ms, object_size, put_path, compressed, encrypted, bucket, key, "Small PUT phase is slow" + duration_ms, object_size, put_path, encrypted, bucket, key, "Small PUT phase is slow" ); } else { debug!( phase, - duration_ms, object_size, put_path, compressed, encrypted, bucket, key, "Small PUT phase exceeded debug threshold" + duration_ms, object_size, put_path, encrypted, bucket, key, "Small PUT phase exceeded debug threshold" ); } } @@ -274,9 +242,8 @@ impl DefaultObjectUsecase { pub(super) async fn run_put_object_flow( input: PutObjectInput, request_context: PutObjectRequestContext, - request_method_name: &'static str, resolved_size: i64, - ) -> S3Result { + ) -> S3Result<(PutObjectOutput, ObjectInfo)> { let start_time = std::time::Instant::now(); let PutObjectInput { @@ -315,7 +282,7 @@ impl DefaultObjectUsecase { let server_side_encryption = server_side_encryption.or(extract_server_side_encryption_from_headers(&request_context.headers)?); - validate_object_key(&key, request_method_name)?; + validate_object_key(&key, if request_context.is_post_object { "POST" } else { "PUT" })?; let Some(body) = body else { return Err(s3_error!(IncompleteBody)) }; @@ -325,33 +292,11 @@ impl DefaultObjectUsecase { let mut small_object_eager_stage = false; let bytes_pool = get_concurrency_manager().bytes_pool(); - let store = get_validated_store_adapter(&bucket).await?; + let store = get_validated_store(&bucket).await?; - let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok(); - - let mut effective_sse = server_side_encryption.or_else(|| { - bucket_sse_config.as_ref().and_then(|(config, _timestamp)| { - config.rules.first().and_then(|rule| { - rule.apply_server_side_encryption_by_default - .as_ref() - .map(|sse| match sse.sse_algorithm.as_str() { - "AES256" => ServerSideEncryption::from_static(ServerSideEncryption::AES256), - "aws:kms" => ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS), - _ => ServerSideEncryption::from_static(ServerSideEncryption::AES256), - }) - }) - }) - }); - - let mut effective_kms_key_id = ssekms_key_id.or_else(|| { - bucket_sse_config.as_ref().and_then(|(config, _timestamp)| { - config.rules.first().and_then(|rule| { - rule.apply_server_side_encryption_by_default - .as_ref() - .and_then(|sse| sse.kms_master_key_id.clone()) - }) - }) - }); + let (default_sse, default_kms_key_id) = resolve_bucket_default_server_side_encryption(&bucket).await; + let mut effective_sse = server_side_encryption.or(default_sse); + let mut effective_kms_key_id = ssekms_key_id.or(default_kms_key_id); validate_sse_headers_for_write( effective_sse.as_ref(), @@ -423,8 +368,12 @@ impl DefaultObjectUsecase { .await?; let eager_max_bytes = resolved_small_put_eager_max_bytes(topology_aware_small_put_eager_max_bytes(&store, opts.versioned)); - let can_use_small_put_eager = - !request_uses_trailing_checksum(&request_context.headers, &request_context.trailing_headers); + let can_use_small_put_eager = request_context.trailing_headers.is_none() + && !request_context.headers.contains_key(AMZ_TRAILER) + && !matches!( + rustfs_rio::get_content_checksum(&request_context.headers), + Ok(Some(checksum)) if checksum.checksum_type.trailing() + ); let current_opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &request_context.headers) .await @@ -539,15 +488,22 @@ impl DefaultObjectUsecase { stage }; + let put_path = if small_object_eager_stage { + "small_eager" + } else if transform_stage.compression_applied() { + "compressed" + } else if plain_reduced_copy_stage { + "reduced_copy" + } else { + "legacy_plain" + }; log_put_flow_phase( &bucket, &key, "build_hash_stage", reader_stage_start.elapsed(), actual_size, - small_object_eager_stage, - plain_reduced_copy_stage, - transform_stage.compression_applied(), + put_path, false, ); let mut reader = stage.reader; @@ -615,19 +571,17 @@ impl DefaultObjectUsecase { "store_put_object", store_put_start.elapsed(), actual_size, - small_object_eager_stage, - plain_reduced_copy_stage, - transform_stage.compression_applied(), + put_path, transform_stage.encryption_applied(), ); - maybe_enqueue_transition_immediate(&obj_info, LcEventSrc::S3PutObject).await; + enqueue_transition_immediate(&obj_info, LcEventSrc::S3PutObject).await; rustfs_ecstore::data_usage::increment_bucket_usage_memory(&bucket, obj_info.size as u64).await; let raw_version = obj_info.version_id.map(|v| v.to_string()); - let put_version = if bucket_prefix_versioning_enabled(&bucket, &key).await { + let put_version = if BucketVersioningSys::prefix_enabled(&bucket, &key).await { raw_version.clone() } else { None @@ -712,11 +666,7 @@ impl DefaultObjectUsecase { } } - Ok(PutObjectFlowResult { - output, - helper_object: obj_info, - helper_version_id: raw_version, - }) + Ok((output, obj_info)) } } diff --git a/rustfs/src/app/object_usecase/types.rs b/rustfs/src/app/object_usecase/types.rs index a0213541d..c7f0288c8 100644 --- a/rustfs/src/app/object_usecase/types.rs +++ b/rustfs/src/app/object_usecase/types.rs @@ -18,12 +18,10 @@ use super::*; pub(super) struct GetObjectRequestContext { pub(super) bucket: String, pub(super) key: String, - pub(super) version_id_for_event: String, pub(super) part_number: Option, pub(super) rs: Option, pub(super) opts: ObjectOptions, pub(super) headers: HeaderMap, - pub(super) method: hyper::Method, pub(super) sse_customer_key: Option, pub(super) sse_customer_key_md5: Option, } @@ -43,9 +41,3 @@ pub(super) struct PutObjectRequestContext { pub(super) region: Option, pub(super) service: Option, } - -pub(super) struct PutObjectFlowResult { - pub(super) output: PutObjectOutput, - pub(super) helper_object: ObjectInfo, - pub(super) helper_version_id: Option, -} diff --git a/rustfs/src/app/object_usecase/zero_copy_tests.rs b/rustfs/src/app/object_usecase/zero_copy_tests.rs index 3aa253bdf..1554dd3d7 100644 --- a/rustfs/src/app/object_usecase/zero_copy_tests.rs +++ b/rustfs/src/app/object_usecase/zero_copy_tests.rs @@ -40,6 +40,58 @@ static DIRECT_CHUNK_TEST_ENV: OnceLock<(Vec, Arc)> = OnceLock: static DIRECT_CHUNK_MULTI_DISK_TEST_ENV: OnceLock<(Vec, Arc)> = OnceLock::new(); static DIRECT_CHUNK_TEST_INIT: Once = Once::new(); +async fn prepare_get_object_request_context(req: &S3Request) -> S3Result { + let GetObjectInput { + bucket, + key, + version_id, + part_number, + range, + .. + } = req.input.clone(); + + validate_object_key(&key, "GET")?; + + let part_number = part_number.map(|value| value as usize); + if let Some(part_number) = part_number + && part_number == 0 + { + return Err(s3_error!(InvalidArgument, "Invalid part number: part number must be greater than 0")); + } + + let rs = range.map(|value| match value { + Range::Int { first, last } => HTTPRangeSpec { + is_suffix_length: false, + start: first as i64, + end: last.map_or(-1, |last| last as i64), + }, + Range::Suffix { length } => HTTPRangeSpec { + is_suffix_length: true, + start: length as i64, + end: -1, + }, + }); + + if rs.is_some() && part_number.is_some() { + return Err(s3_error!(InvalidArgument, "range and part_number invalid")); + } + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), part_number, &req.headers) + .await + .map_err(ApiError::from)?; + + Ok(GetObjectRequestContext { + bucket, + key, + part_number, + rs, + opts, + headers: req.headers.clone(), + sse_customer_key: req.input.sse_customer_key.clone(), + sse_customer_key_md5: req.input.sse_customer_key_md5.clone(), + }) +} + fn init_direct_chunk_test_tracing() { DIRECT_CHUNK_TEST_INIT.call_once(|| {}); } @@ -287,20 +339,11 @@ async fn select_reconstructed_chunk_read( continue; } - let candidate = get_object_zero_copy::prepare_get_object_chunk_read( - request_context, - ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let candidate = + get_object_zero_copy::prepare_get_object_chunk_read(request_context, ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); let is_reconstructed = matches!( &candidate.body_source, @@ -377,20 +420,11 @@ async fn prepare_get_object_chunk_read_marks_direct_path_for_single_disk_store() let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); match read_setup.body_source { GetObjectBodySource::Chunk { path, copy_mode, .. } => { @@ -441,19 +475,10 @@ async fn prepare_get_object_chunk_read_falls_back_to_legacy_when_chunk_bridge_fa let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap(); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap(); assert!(read_setup.is_none(), "chunk bridge failure should fall back to legacy reader"); } @@ -489,20 +514,11 @@ async fn execute_get_object_range_marks_direct_path_for_single_disk_store() { let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); match read_setup.body_source { GetObjectBodySource::Chunk { path, .. } => { @@ -562,20 +578,11 @@ async fn execute_get_object_range_marks_direct_path_for_multi_disk_store_without let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); match read_setup.body_source { GetObjectBodySource::Chunk { path, copy_mode, .. } => { @@ -646,20 +653,11 @@ async fn execute_get_object_range_marks_reconstructed_path_for_multi_disk_store_ let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); match read_setup.body_source { GetObjectBodySource::Chunk { path, copy_mode, .. } => { @@ -708,20 +706,11 @@ async fn execute_get_object_part_number_marks_direct_path_for_single_disk_store( let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); match read_setup.body_source { GetObjectBodySource::Chunk { path, copy_mode, .. } => { @@ -765,20 +754,11 @@ async fn execute_get_object_whole_multipart_marks_direct_path_for_single_disk_st let request_context = prepare_get_object_request_context(&req).await.unwrap(); let manager = get_concurrency_manager(); - let read_setup = get_object_zero_copy::prepare_get_object_chunk_read( - &request_context, - &ecstore, - manager, - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - request_context.part_number, - &request_context.opts, - std::time::Instant::now(), - ) - .await - .unwrap() - .expect("expected chunk fast path"); + let read_setup = + get_object_zero_copy::prepare_get_object_chunk_read(&request_context, &ecstore, manager, std::time::Instant::now()) + .await + .unwrap() + .expect("expected chunk fast path"); match read_setup.body_source { GetObjectBodySource::Chunk { path, copy_mode, .. } => {