From df109680ae5744e56068bfe8c33209ed602567cc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Wed, 15 Apr 2026 11:50:32 +0800 Subject: [PATCH] cleanup: remove orphaned chunk fast-path submodule files (#2545) Co-authored-by: cxymds --- .../src/app/object_usecase/get_object_flow.rs | 457 ---------- .../object_usecase/get_object_zero_copy.rs | 230 ----- .../app/object_usecase/put_object_extract.rs | 514 ----------- .../src/app/object_usecase/put_object_flow.rs | 798 ------------------ 4 files changed, 1999 deletions(-) delete mode 100644 rustfs/src/app/object_usecase/get_object_flow.rs delete mode 100644 rustfs/src/app/object_usecase/get_object_zero_copy.rs delete mode 100644 rustfs/src/app/object_usecase/put_object_extract.rs delete mode 100644 rustfs/src/app/object_usecase/put_object_flow.rs diff --git a/rustfs/src/app/object_usecase/get_object_flow.rs b/rustfs/src/app/object_usecase/get_object_flow.rs deleted file mode 100644 index 53c6497f0..000000000 --- a/rustfs/src/app/object_usecase/get_object_flow.rs +++ /dev/null @@ -1,457 +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::DeadlockRequestGuard; -use super::GetObjectRequestContext; -use super::get_object_zero_copy::{GetObjectIoPlanning, GetObjectPreparedRead, prepare_get_object_read_execution}; -use crate::error::ApiError; -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::{ - GetObjectBodyPlan as ObjectIoGetObjectBodyPlan, GetObjectBodyPlanningInputs as ObjectIoGetObjectBodyPlanningInputs, - GetObjectBodySource, GetObjectDataPlaneMetricContract as ObjectIoGetObjectDataPlaneMetricContract, GetObjectFlowResult, - GetObjectOutputContext, GetObjectReadSetup, MaterializeGetObjectBodyError as ObjectIoMaterializeGetObjectBodyError, - 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, - 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::StreamingBlob; -use std::time::Duration; -use tokio::io::AsyncRead; -use tracing::{debug, error, info, warn}; - -pub(super) struct GetObjectBootstrap { - pub(super) timeout_config: TimeoutConfig, - pub(super) wrapper: RequestTimeoutWrapper, - pub(super) request_start: std::time::Instant, - pub(super) request_guard: GetObjectGuard, - pub(super) _deadlock_request_guard: DeadlockRequestGuard, -} - -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 -} - -pub(super) async fn build_get_object_output_context( - request_context: &GetObjectRequestContext, - manager: &ConcurrencyManager, - read_setup: GetObjectReadSetup, - io_planning: &GetObjectIoPlanning<'_>, - base_buffer_size: usize, - versioned: bool, -) -> S3Result<(GetObjectOutputContext, ObjectIoGetObjectDataPlaneMetricContract)> { - 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, - 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 GetObjectBodySource::Reader(final_stream) = body_source; - - let body = build_get_object_body_adapter( - final_stream, - bucket, - key, - response_content_length, - optimal_buffer_size, - ObjectIoGetObjectBodyPlanningInputs { - is_part_request: part_number.is_some(), - is_range_request: rs.is_some(), - encryption_applied, - response_size: response_content_length, - }, - ) - .await?; - let metric_contract = ObjectIoGetObjectDataPlaneMetricContract::disk( - rustfs_io_metrics::IoPath::Legacy, - rustfs_io_metrics::CopyMode::SingleCopy, - ); - - let checksums = object_io_build_get_object_checksums(&info, &request_context.headers, part_number, rs.as_ref()) - .map_err(ApiError::from)?; - let filtered_metadata = filter_object_metadata(&info.user_defined); - - Ok(( - object_io_build_get_object_output_context( - body, - info, - event_info, - content_type, - last_modified, - response_content_length, - content_range, - server_side_encryption, - sse_customer_algorithm, - sse_customer_key_md5, - ssekms_key_id, - &checksums, - filtered_metadata, - versioned, - optimal_buffer_size, - Some(metric_contract.copy_mode), - ), - metric_contract, - )) -} - -pub(super) async fn run_get_object_flow( - request_context: GetObjectRequestContext, - version_id_for_event: String, - manager: &ConcurrencyManager, - bootstrap: &GetObjectBootstrap, - base_buffer_size: usize, -) -> S3Result { - let timeout_config = &bootstrap.timeout_config; - let wrapper = &bootstrap.wrapper; - let request_start = bootstrap.request_start; - - let prepared_read = prepare_get_object_read_execution(&request_context, manager, wrapper, timeout_config).await?; - let GetObjectPreparedRead { io_planning, read_setup } = prepared_read; - - 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( - &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 futures_util::StreamExt; - use http::HeaderMap; - use rustfs_ecstore::store_api::ObjectOptions; - use rustfs_object_io::get::{GetObjectEncryptionState, LegacyReadPlan, build_reader_read_setup}; - use rustfs_rio::{Reader, WarpReader}; - use std::{io::Cursor, sync::Arc, time::Duration}; - use tokio::sync::Semaphore; - - 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, - } - } - - #[tokio::test] - async fn build_get_object_output_context_materializes_reader_payload_with_legacy_metrics() { - let payload = b"hello from legacy reader".to_vec(); - let reader = Box::new(WarpReader::new(Cursor::new(payload.clone()))) as Box; - let read_setup = build_reader_read_setup( - ObjectInfo::default(), - ObjectInfo::default(), - reader, - LegacyReadPlan { - rs: None, - content_type: None, - last_modified: None, - response_content_length: payload.len() as i64, - content_range: None, - }, - GetObjectEncryptionState::default(), - ); - let semaphore = Arc::new(Semaphore::new(1)); - let permit = semaphore.acquire().await.expect("disk permit"); - let io_planning = GetObjectIoPlanning { - _disk_permit: permit, - permit_wait_duration: Duration::ZERO, - queue_status: crate::storage::concurrency::IoQueueStatus { - total_permits: 1, - permits_in_use: 1, - ..Default::default() - }, - queue_utilization: 100.0, - }; - let manager = ConcurrencyManager::new(); - - let (output_context, metric_contract) = - build_get_object_output_context(&sample_request_context(), &manager, read_setup, &io_planning, 8 * 1024, false) - .await - .expect("reader-backed output context"); - - assert_eq!(metric_contract.io_path, rustfs_io_metrics::IoPath::Legacy); - assert_eq!(metric_contract.copy_mode, rustfs_io_metrics::CopyMode::SingleCopy); - assert_eq!(output_context.output.content_length, Some(payload.len() as i64)); - assert_eq!(output_context.copy_mode_override, Some(rustfs_io_metrics::CopyMode::SingleCopy)); - - let mut body = output_context.output.body.expect("streaming body"); - let mut collected = Vec::new(); - while let Some(chunk) = body.next().await { - collected.extend_from_slice(&chunk.expect("body chunk")); - } - assert_eq!(collected, payload); - } - - #[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 deleted file mode 100644 index 42f9e5b74..000000000 --- a/rustfs/src/app/object_usecase/get_object_zero_copy.rs +++ /dev/null @@ -1,230 +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::GetObjectRequestContext; -use crate::error::ApiError; -use crate::storage::concurrency::{self, ConcurrencyManager}; -use crate::storage::timeout_wrapper::{RequestTimeoutWrapper, TimeoutConfig}; -use crate::storage::{ - 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::ObjectIO; -use rustfs_object_io::get::{ - GetObjectEncryptionState as ObjectIoGetObjectEncryptionState, GetObjectReadSetup, - build_reader_read_setup as object_io_build_reader_read_setup, plan_legacy_read as object_io_plan_legacy_read, -}; -use rustfs_rio::{Reader, WarpReader}; -use s3s::{S3Result, s3_error}; -use std::time::Duration; -use tracing::{debug, warn}; - -pub(super) struct GetObjectIoPlanning<'a> { - pub(super) _disk_permit: tokio::sync::SemaphorePermit<'a>, - pub(super) permit_wait_duration: Duration, - pub(super) queue_status: concurrency::IoQueueStatus, - pub(super) queue_utilization: f64, -} - -pub(super) struct GetObjectPreparedRead<'a> { - pub(super) io_planning: GetObjectIoPlanning<'a>, - pub(super) read_setup: GetObjectReadSetup, -} - -pub(super) async fn acquire_get_object_io_planning<'a>( - manager: &'a ConcurrencyManager, - wrapper: &RequestTimeoutWrapper, - timeout_config: &TimeoutConfig, - bucket: &str, - key: &str, -) -> S3Result> { - let permit_wait_start = std::time::Instant::now(); - let disk_permit = manager - .acquire_disk_read_permit() - .await - .map_err(|_| s3_error!(InternalError, "disk read semaphore closed"))?; - let permit_wait_duration = permit_wait_start.elapsed(); - - if wrapper.is_timeout() { - warn!( - bucket = %bucket, - key = %key, - wait_ms = permit_wait_duration.as_millis(), - timeout_secs = timeout_config.get_object_timeout.as_secs(), - elapsed_ms = wrapper.elapsed().as_millis(), - "GetObject request timed out while waiting for disk permit" - ); - - rustfs_io_metrics::record_get_object_timeout(Some("disk_permit"), Some(wrapper.elapsed().as_secs_f64())); - return Err(s3_error!(InternalError, "Request timeout while waiting for disk permit")); - } - - let queue_status = manager.io_queue_status(); - let queue_snapshot = GetObjectQueueSnapshot::from_available_permits( - queue_status.total_permits, - queue_status.total_permits.saturating_sub(queue_status.permits_in_use), - ); - let queue_utilization = queue_snapshot.utilization_percent(); - - if queue_snapshot.is_congested(80.0) { - warn!( - bucket = %bucket, - key = %key, - queue_utilization = format!("{:.1}%", queue_utilization), - permits_in_use = queue_status.permits_in_use, - total_permits = queue_status.total_permits, - "I/O queue congestion detected" - ); - - rustfs_io_metrics::record_io_queue_congestion(); - } - - 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 reading object" - ); - rustfs_io_metrics::record_get_object_timeout(Some("before_read"), Some(wrapper.elapsed().as_secs_f64())); - return Err(s3_error!(InternalError, "Request timeout before reading object")); - } - - Ok(GetObjectIoPlanning { - _disk_permit: disk_permit, - permit_wait_duration, - queue_status, - queue_utilization, - }) -} - -pub(super) async fn prepare_get_object_read( - request_context: &GetObjectRequestContext, - store: &rustfs_ecstore::store::ECStore, - manager: &ConcurrencyManager, - read_start: std::time::Instant, -) -> S3Result { - let reader = store - .get_object_reader( - &request_context.bucket, - &request_context.key, - request_context.rs.clone(), - HeaderMap::new(), - &request_context.opts, - ) - .await - .map_err(ApiError::from)?; - - let info = reader.object_info; - - let read_duration = read_start.elapsed(); - rustfs_io_metrics::record_io_path_selected("get", rustfs_io_metrics::IoPath::Legacy); - - manager.record_disk_operation(info.size as u64, read_duration, true).await; - - check_preconditions(&request_context.headers, &info)?; - - debug!(object_size = info.size, part_count = info.parts.len(), "GET object metadata snapshot"); - for part in &info.parts { - debug!( - part_number = part.number, - part_size = part.size, - part_actual_size = part.actual_size, - "GET object part details" - ); - } - - let event_info = info.clone(); - validate_sse_headers_for_read(&info.user_defined, &request_context.headers)?; - validate_ssec_for_read( - &info.user_defined, - request_context.sse_customer_key.as_ref(), - request_context.sse_customer_key_md5.as_ref(), - )?; - 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={:?}", - info.parts.len(), - request_context.sse_customer_key.is_some() - ); - - let decryption_request = DecryptionRequest { - 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(), - part_number: None, - parts: &info.parts, - etag: info.etag.as_deref(), - }; - - let encrypted_stream = reader.stream; - - let (encryption_state, final_stream) = match sse_decryption(decryption_request).await? { - Some(material) => { - let server_side_encryption = Some(material.server_side_encryption.clone()); - let sse_customer_algorithm = Some(material.algorithm.clone()); - let sse_customer_key_md5 = material.customer_key_md5.clone(); - let ssekms_key_id = material.kms_key_id.clone(); - let (decrypted_stream, plaintext_size) = material - .wrap_reader(encrypted_stream, read_plan.response_content_length) - .await - .map_err(ApiError::from)?; - - ( - ObjectIoGetObjectEncryptionState { - server_side_encryption, - sse_customer_algorithm, - sse_customer_key_md5, - ssekms_key_id, - encryption_applied: true, - response_content_length_override: Some(plaintext_size), - }, - decrypted_stream, - ) - } - None => ( - ObjectIoGetObjectEncryptionState::default(), - Box::new(WarpReader::new(encrypted_stream)) as Box, - ), - }; - - Ok(object_io_build_reader_read_setup( - info, - event_info, - final_stream, - read_plan, - encryption_state, - )) -} - -pub(super) async fn prepare_get_object_read_execution<'a>( - request_context: &GetObjectRequestContext, - manager: &'a ConcurrencyManager, - wrapper: &RequestTimeoutWrapper, - timeout_config: &TimeoutConfig, -) -> S3Result> { - 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_setup = prepare_get_object_read(request_context, &store, manager, std::time::Instant::now()).await?; - - Ok(GetObjectPreparedRead { io_planning, read_setup }) -} diff --git a/rustfs/src/app/object_usecase/put_object_extract.rs b/rustfs/src/app/object_usecase/put_object_extract.rs deleted file mode 100644 index 46caf7a75..000000000 --- a/rustfs/src/app/object_usecase/put_object_extract.rs +++ /dev/null @@ -1,514 +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::*; -use crate::app::context::NotifyInterface; -use rustfs_object_io::put::{ - apply_extract_entry_pax_extensions, apply_trailing_checksums, is_sse_kms_requested, map_extract_archive_error, - normalize_extract_entry_key, resolve_put_object_extract_options, -}; - -impl DefaultObjectUsecase { - pub(super) async fn run_put_object_extract_flow( - input: PutObjectInput, - request_context: PutObjectRequestContext, - notify: Arc, - resolved_size: i64, - ) -> S3Result { - if is_sse_kms_requested(&input, &request_context.headers) { - return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for extract uploads")); - } - - let PutObjectInput { - body, - bucket, - key, - version_id, - cache_control, - content_disposition, - content_encoding, - content_length: _content_length, - content_language, - content_type, - content_md5, - expires, - object_lock_legal_hold_status, - object_lock_mode, - object_lock_retain_until_date, - server_side_encryption, - sse_customer_algorithm, - sse_customer_key, - sse_customer_key_md5, - ssekms_key_id, - storage_class, - tagging, - website_redirect_location, - .. - } = input; - - let event_version_id = version_id; - let (h_algo, h_key, h_md5) = extract_ssec_params_from_headers(&request_context.headers)?; - let sse_customer_algorithm = sse_customer_algorithm.or(h_algo); - let sse_customer_key = sse_customer_key.or(h_key); - 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 (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)) - { - return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for extract uploads")); - } - validate_sse_headers_for_write( - effective_sse.as_ref(), - effective_kms_key_id.as_ref(), - sse_customer_algorithm.as_ref(), - sse_customer_key.as_ref(), - sse_customer_key_md5.as_ref(), - true, - )?; - let Some(body) = body else { return Err(s3_error!(IncompleteBody)) }; - - let size = resolved_size; - validate_object_key(&key, "PUT")?; - - let buffer_size = get_buffer_size_opt_in(size); - let body = tokio::io::BufReader::with_capacity( - buffer_size, - StreamReader::new(body.map(|f| f.map_err(|e| std::io::Error::other(e.to_string())))), - ); - - let Some(ext) = Path::new(&key).extension().and_then(|s| s.to_str()) else { - return Err(s3_error!(InvalidArgument, "key extension not found")); - }; - - let ext = ext.to_owned(); - - let md5hex = if let Some(base64_md5) = content_md5 { - let md5 = base64_simd::STANDARD - .decode_to_vec(base64_md5.as_bytes()) - .map_err(|e| ApiError::from(StorageError::other(format!("Invalid content MD5: {e}"))))?; - Some(hex_simd::encode_to_string(&md5, hex_simd::AsciiCase::Lower)) - } else { - None - }; - - let sha256hex = get_content_sha256_with_query(&request_context.headers, request_context.uri_query.as_deref()); - let actual_size = size; - - let mut archive_reader = - HashReader::from_stream(body, size, actual_size, md5hex, sha256hex, false).map_err(ApiError::from)?; - - if let Err(err) = - archive_reader.add_checksum_from_s3s(&request_context.headers, request_context.trailing_headers.clone(), false) - { - return Err(ApiError::from(err).into()); - } - - let archive_etag = Arc::new(Mutex::new(None)); - let decoder = CompressionFormat::from_extension(&ext) - .get_decoder(ExtractArchiveEtagReader::new(archive_reader, archive_etag.clone())) - .map_err(|e| { - error!("get_decoder err {:?}", e); - s3_error!(InvalidArgument, "get_decoder err") - })?; - - let mut ar = Archive::new(decoder); - let mut entries = ar.entries().map_err(|e| { - error!("get entries err {:?}", e); - s3_error!(InvalidArgument, "get entries err") - })?; - - let Some(store) = new_object_layer_fn() else { - return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); - }; - - let extract_options = resolve_put_object_extract_options(&request_context.headers); - let version_id = match event_version_id { - Some(v) => v.to_string(), - None => String::new(), - }; - - let req_params = extract_params_header(&request_context.headers); - let host = get_request_host(&request_context.headers); - let port = get_request_port(&request_context.headers); - let user_agent = get_request_user_agent(&request_context.headers); - let tracing_context = request_context.extensions.get::().cloned(); - - while let Some(entry) = entries.next().await { - let mut f = match entry { - Ok(f) => f, - Err(e) => { - if extract_options.ignore_errors { - warn!("Skipping archive entry because read failed and ignore-errors is enabled: {e}"); - continue; - } - error!("Failed to read archive entry: {}", e); - return Err(s3_error!(InvalidArgument, "Failed to read archive entry: {:?}", e)); - } - }; - - let fpath = match f.path() { - Ok(path) => path, - Err(e) => { - if extract_options.ignore_errors { - warn!("Skipping archive entry because path decode failed and ignore-errors is enabled: {e}"); - continue; - } - return Err(s3_error!(InvalidArgument, "Failed to decode archive entry path")); - } - }; - - 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); - - 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 - .header() - .mtime() - .ok() - .and_then(|modified_at_secs| OffsetDateTime::from_unix_timestamp(modified_at_secs as i64).ok()); - let mut metadata = HashMap::new(); - apply_put_request_metadata( - &mut metadata, - &request_context.headers, - &fpath, - cache_control.clone(), - content_disposition.clone(), - content_encoding.clone(), - content_language.clone(), - content_type.clone(), - expires.clone(), - website_redirect_location.clone(), - tagging.clone(), - storage_class.clone(), - )?; - let mut opts = put_opts(&bucket, &fpath, None, &request_context.headers, metadata.clone()) - .await - .map_err(ApiError::from)?; - apply_extract_entry_pax_extensions(&mut f, &mut metadata, &mut opts).await?; - if archive_entry_mod_time.is_some() { - opts.mod_time = archive_entry_mod_time; - } - - debug!("Extracting file: {}, size: {} bytes", fpath, size); - - if is_dir { - if extract_options.ignore_dirs { - debug!("Skipping directory entry during archive extract: {}", fpath); - continue; - } - size = 0; - } - - let actual_size = size; - let should_compress = !is_dir && is_compressible(&HeaderMap::new(), &fpath) && size > MIN_COMPRESSIBLE_SIZE as i64; - - let mut hrd = if is_dir { - HashReader::from_stream(std::io::Cursor::new(Vec::new()), size, actual_size, None, None, false) - .map_err(ApiError::from)? - } else if should_compress { - insert_str(&mut metadata, SUFFIX_COMPRESSION, CompressionAlgorithm::default().to_string()); - insert_str(&mut metadata, SUFFIX_ACTUAL_SIZE, size.to_string()); - - let hrd = HashReader::from_stream(f, size, actual_size, None, None, false).map_err(ApiError::from)?; - size = HashReader::SIZE_PRESERVE_LAYER; - HashReader::from_reader( - CompressReader::new(hrd, CompressionAlgorithm::default()), - size, - actual_size, - None, - None, - false, - ) - .map_err(ApiError::from)? - } else { - HashReader::from_stream(f, size, actual_size, None, None, false).map_err(ApiError::from)? - }; - apply_put_request_object_lock_opts( - &bucket, - object_lock_legal_hold_status.clone(), - object_lock_mode.clone(), - object_lock_retain_until_date.clone(), - &mut opts, - ) - .await?; - if let Some(material) = sse_encryption(EncryptionRequest { - bucket: &bucket, - key: &fpath, - server_side_encryption: effective_sse.clone(), - ssekms_key_id: effective_kms_key_id.clone(), - sse_customer_algorithm: sse_customer_algorithm.clone(), - sse_customer_key: sse_customer_key.clone(), - sse_customer_key_md5: sse_customer_key_md5.clone(), - content_size: actual_size, - part_number: None, - part_key: None, - part_nonce: None, - }) - .await? - { - effective_sse = Some(material.server_side_encryption.clone()); - effective_kms_key_id = material.kms_key_id.clone(); - - let encrypted_reader = material.wrap_reader(hrd); - hrd = HashReader::from_reader(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false) - .map_err(ApiError::from)?; - - let encryption_metadata = material.metadata; - metadata.extend(encryption_metadata.clone()); - opts.user_defined.extend(encryption_metadata); - } - opts.user_defined.extend(metadata); - let capacity_scope_token = Uuid::new_v4(); - opts.capacity_scope_token = Some(capacity_scope_token); - let mut reader = PutObjReader::new(hrd); - - let obj_info = match store.put_object(&bucket, &fpath, &mut reader, &opts).await { - Ok(info) => info, - Err(e) => { - if extract_options.ignore_errors { - warn!("Skipping archive entry because object write failed and ignore-errors is enabled: {e}"); - continue; - } - return Err(ApiError::from(e).into()); - } - }; - record_capacity_write(Some(capacity_scope_token)).await; - - let e_tag = obj_info.etag.clone().map(|etag| to_s3s_etag(&etag)); - - let output = PutObjectOutput { - e_tag, - ..Default::default() - }; - - 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: 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 { - crc32: input.checksum_crc32, - crc32c: input.checksum_crc32c, - sha1: input.checksum_sha1, - sha256: input.checksum_sha256, - crc64nvme: input.checksum_crc64nvme, - }; - apply_trailing_checksums( - input.checksum_algorithm.as_ref().map(|a| a.as_str()), - &request_context.trailing_headers, - &mut checksums, - ); - - drop(entries); - let mut decoder = match ar.into_inner() { - Ok(decoder) => decoder, - Err(_) => return Err(s3_error!(InvalidArgument, "Failed to finalize archive reader")), - }; - tokio::io::copy(&mut decoder, &mut tokio::io::sink()) - .await - .map_err(map_extract_archive_error)?; - let archive_etag = archive_etag - .lock() - .ok() - .and_then(|etag| etag.clone()) - .map(|etag| to_s3s_etag(&etag)); - - let output = PutObjectOutput { - e_tag: archive_etag, - checksum_crc32: checksums.crc32, - checksum_crc32c: checksums.crc32c, - checksum_sha1: checksums.sha1, - checksum_sha256: checksums.sha256, - checksum_crc64nvme: checksums.crc64nvme, - ..Default::default() - }; - Ok(output) - } -} - -#[cfg(test)] -mod tests { - use super::*; - use http::{Extensions, HeaderMap, HeaderValue, Method, Uri}; - use rustfs_utils::http::headers::{AMZ_SERVER_SIDE_ENCRYPTION, AMZ_SERVER_SIDE_ENCRYPTION_KMS_ID, AMZ_SNOWBALL_EXTRACT}; - - fn build_request(input: T, method: Method) -> S3Request { - S3Request { - input, - method, - uri: Uri::from_static("/"), - headers: HeaderMap::new(), - extensions: Extensions::new(), - credentials: None, - region: None, - service: None, - trailing_headers: None, - } - } - - #[tokio::test] - async fn execute_put_object_rejects_post_object_sse_kms_from_input() { - let input = PutObjectInput::builder() - .bucket("test-bucket".to_string()) - .key("test-key".to_string()) - .server_side_encryption(Some(ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS))) - .build() - .unwrap(); - - let mut req = build_request(input, Method::POST); - req.extensions.insert(PostObjectRequestMarker); - - let usecase = DefaultObjectUsecase::without_context(); - let fs = FS::new(); - - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); - } - - #[tokio::test] - async fn execute_put_object_rejects_extract_sse_kms() { - let input = PutObjectInput::builder() - .bucket("test-bucket".to_string()) - .key("archive.tar".to_string()) - .server_side_encryption(Some(ServerSideEncryption::from_static(ServerSideEncryption::AWS_KMS))) - .build() - .unwrap(); - - let mut req = build_request(input, Method::PUT); - req.headers.insert(AMZ_SNOWBALL_EXTRACT, HeaderValue::from_static("true")); - - let usecase = DefaultObjectUsecase::without_context(); - let fs = FS::new(); - - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); - } - - #[tokio::test] - async fn execute_put_object_extract_rejects_invalid_storage_class() { - let input = PutObjectInput::builder() - .bucket("test-bucket".to_string()) - .key("archive.tar".to_string()) - .storage_class(Some(StorageClass::from_static("INVALID"))) - .build() - .unwrap(); - - let mut req = build_request(input, Method::PUT); - req.headers.insert(AMZ_SNOWBALL_EXTRACT, HeaderValue::from_static("true")); - - let usecase = DefaultObjectUsecase::without_context(); - let fs = FS::new(); - - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); - assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass); - } - - #[tokio::test] - async fn execute_put_object_rejects_post_object_sse_kms_from_headers() { - 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); - req.headers - .insert(AMZ_SERVER_SIDE_ENCRYPTION, HeaderValue::from_static("aws:kms")); - - let usecase = DefaultObjectUsecase::without_context(); - let fs = FS::new(); - - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); - } - - #[tokio::test] - async fn execute_put_object_rejects_post_object_sse_kms_key_id_header() { - 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); - req.headers - .insert(AMZ_SERVER_SIDE_ENCRYPTION_KMS_ID, HeaderValue::from_static("test-kms-key-id")); - - let usecase = DefaultObjectUsecase::without_context(); - let fs = FS::new(); - - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); - } - - #[tokio::test] - async fn execute_put_object_rejects_extract_sse_kms_key_id_header() { - let input = PutObjectInput::builder() - .bucket("test-bucket".to_string()) - .key("archive.tar".to_string()) - .build() - .unwrap(); - - let mut req = build_request(input, Method::PUT); - req.headers.insert(AMZ_SNOWBALL_EXTRACT, HeaderValue::from_static("true")); - req.headers - .insert(AMZ_SERVER_SIDE_ENCRYPTION_KMS_ID, HeaderValue::from_static("test-kms-key-id")); - - let usecase = DefaultObjectUsecase::without_context(); - let fs = FS::new(); - - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); - assert_eq!(err.code(), &S3ErrorCode::NotImplemented); - } -} diff --git a/rustfs/src/app/object_usecase/put_object_flow.rs b/rustfs/src/app/object_usecase/put_object_flow.rs deleted file mode 100644 index 105685b4a..000000000 --- a/rustfs/src/app/object_usecase/put_object_flow.rs +++ /dev/null @@ -1,798 +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::*; -use bytes::Buf; -use futures::{Stream, StreamExt}; -use rustfs_ecstore::config::GLOBAL_STORAGE_CLASS; -use rustfs_io_core::{BytesPool, PooledBuffer}; -use rustfs_object_io::put::{ - PutObjectChecksums, PutObjectIngressKind, PutObjectLegacyHashStagePlan, PutObjectLegacyHashValues, PutObjectTransformStage, - apply_trailing_checksums, build_put_object_ingress_source, build_put_object_legacy_hash_stage, - build_put_object_plain_hash_stage, plan_put_object_body_with_transforms, resolve_put_transformed_fallback_reason, -}; -use rustfs_rio::{EtagResolvable, HashReaderDetector, TryGetIndex}; -use rustfs_utils::http::headers::AMZ_TRAILER; - -const DEFAULT_SMALL_PUT_EAGER_MAX_BYTES: i64 = 1024 * 1024; -const ENV_RUSTFS_PUT_SMALL_EAGER_MAX_BYTES: &str = "RUSTFS_PUT_SMALL_EAGER_MAX_BYTES"; -const ENV_RUSTFS_PUT_FORCE_DISABLE_SMALL_EAGER: &str = "RUSTFS_PUT_FORCE_DISABLE_SMALL_EAGER"; -const SLOW_PUT_PHASE_DEBUG_THRESHOLD_MS: u64 = 100; -const SLOW_PUT_PHASE_WARN_THRESHOLD_MS: u64 = 1_000; -const SLOW_PUT_PHASE_ERROR_THRESHOLD_MS: u64 = 5_000; - -fn resolved_checksum_bytes(checksums: &PutObjectChecksums) -> Option { - [ - (rustfs_rio::ChecksumType::CRC32, checksums.crc32.as_deref()), - (rustfs_rio::ChecksumType::CRC32C, checksums.crc32c.as_deref()), - (rustfs_rio::ChecksumType::SHA1, checksums.sha1.as_deref()), - (rustfs_rio::ChecksumType::SHA256, checksums.sha256.as_deref()), - (rustfs_rio::ChecksumType::CRC64_NVME, checksums.crc64nvme.as_deref()), - ] - .into_iter() - .find_map(|(checksum_type, value)| { - value.and_then(|value| rustfs_rio::Checksum::new_with_type(checksum_type, value).map(|checksum| checksum.to_bytes(&[]))) - }) -} - -fn clamp_small_put_eager_max_bytes(inline_object_limit_bytes: Option) -> i64 { - inline_object_limit_bytes - .unwrap_or(DEFAULT_SMALL_PUT_EAGER_MAX_BYTES as usize) - .min(DEFAULT_SMALL_PUT_EAGER_MAX_BYTES as usize) as i64 -} - -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; - }; - - let data_shards = first_pool - .set_drive_count - .saturating_sub(first_pool.default_parity_count) - .max(1); - - let inline_object_limit = GLOBAL_STORAGE_CLASS - .get() - .map(|config| config.inline_object_limit_bytes(data_shards, versioned)); - - clamp_small_put_eager_max_bytes(inline_object_limit) -} - -fn resolved_small_put_eager_max_bytes(default_max_bytes: i64) -> i64 { - if rustfs_utils::get_env_bool(ENV_RUSTFS_PUT_FORCE_DISABLE_SMALL_EAGER, false) { - return 0; - } - - 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) -} - -fn should_use_small_put_eager_path(size: i64, eager_max_bytes: i64, compression_enabled: bool, encryption_enabled: bool) -> bool { - size > 0 && size <= eager_max_bytes && !compression_enabled && !encryption_enabled -} - -fn log_put_flow_phase( - bucket: &str, - key: &str, - phase: &str, - elapsed: Duration, - object_size: i64, - put_path: &'static str, - encrypted: bool, -) { - let duration_ms = elapsed.as_millis() as u64; - if duration_ms < SLOW_PUT_PHASE_DEBUG_THRESHOLD_MS { - return; - } - - if duration_ms >= SLOW_PUT_PHASE_ERROR_THRESHOLD_MS { - error!( - phase, - 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, encrypted, bucket, key, "Small PUT phase is slow" - ); - } else { - debug!( - phase, - duration_ms, object_size, put_path, encrypted, bucket, key, "Small PUT phase exceeded debug threshold" - ); - } -} - -struct PooledBufferReader { - buffer: PooledBuffer, - position: usize, -} - -impl PooledBufferReader { - fn new(buffer: PooledBuffer) -> Self { - Self { buffer, position: 0 } - } -} - -impl AsyncRead for PooledBufferReader { - fn poll_read( - mut self: std::pin::Pin<&mut Self>, - _cx: &mut std::task::Context<'_>, - buf: &mut ReadBuf<'_>, - ) -> std::task::Poll> { - let remaining = &self.buffer[self.position..]; - if remaining.is_empty() { - return std::task::Poll::Ready(Ok(())); - } - - let to_copy = remaining.len().min(buf.remaining()); - buf.put_slice(&remaining[..to_copy]); - self.position += to_copy; - std::task::Poll::Ready(Ok(())) - } -} - -impl EtagResolvable for PooledBufferReader {} - -impl HashReaderDetector for PooledBufferReader {} - -impl TryGetIndex for PooledBufferReader {} - -async fn read_small_put_body_eager(body: S, size: i64, pool: Arc) -> S3Result -where - S: Stream>, - B: Buf, - E: std::fmt::Display, -{ - let expected_len = usize::try_from(size).map_err(|_| s3_error!(InvalidRequest, "Object size overflow"))?; - let mut data = pool.acquire_buffer(expected_len).await; - let mut body = Box::pin(body); - - while let Some(result) = body.next().await { - let mut chunk = result.map_err(|err| S3Error::with_message(S3ErrorCode::IncompleteBody, err.to_string()))?; - let chunk_len = chunk.remaining(); - if chunk_len == 0 { - continue; - } - - let new_len = data - .len() - .checked_add(chunk_len) - .ok_or_else(|| s3_error!(InvalidRequest, "Object size overflow"))?; - if new_len > expected_len { - return Err(s3_error!(IncompleteBody)); - } - - let start = data.len(); - data.resize(new_len, 0); - chunk.copy_to_slice(&mut data[start..new_len]); - - if data.len() == expected_len { - return Ok(data); - } - } - - if data.len() != expected_len { - return Err(s3_error!(IncompleteBody)); - } - - Ok(data) -} - -async fn build_small_put_eager_hash_stage( - body: S, - size: i64, - pool: Arc, - hash_values: PutObjectLegacyHashValues, - headers: &HeaderMap, - trailing_headers: Option, -) -> S3Result -where - S: Stream>, - B: Buf, - E: std::fmt::Display, -{ - let data = read_small_put_body_eager(body, size, pool).await?; - build_put_object_legacy_hash_stage( - Box::new(PooledBufferReader::new(data)), - hash_values, - PutObjectLegacyHashStagePlan { - size, - actual_size: size, - apply_s3_checksum: true, - ignore_s3_checksum_value: false, - }, - headers, - trailing_headers, - ) - .map_err(ApiError::from) - .map_err(Into::into) -} - -impl DefaultObjectUsecase { - pub(super) async fn run_put_object_flow( - input: PutObjectInput, - request_context: PutObjectRequestContext, - resolved_size: i64, - ) -> S3Result<(PutObjectOutput, ObjectInfo)> { - let start_time = std::time::Instant::now(); - - let PutObjectInput { - body, - bucket, - cache_control, - key, - content_length: _content_length, - content_disposition, - content_encoding, - content_language, - content_type, - expires, - tagging, - metadata, - version_id, - server_side_encryption, - sse_customer_algorithm, - sse_customer_key, - sse_customer_key_md5, - ssekms_key_id, - content_md5, - object_lock_legal_hold_status, - object_lock_mode, - object_lock_retain_until_date, - storage_class, - website_redirect_location, - .. - } = input; - - let (h_algo, h_key, h_md5) = extract_ssec_params_from_headers(&request_context.headers)?; - let sse_customer_algorithm = sse_customer_algorithm.or(h_algo); - let sse_customer_key = sse_customer_key.or(h_key); - let sse_customer_key_md5 = sse_customer_key_md5.or(h_md5); - - let server_side_encryption = - server_side_encryption.or(extract_server_side_encryption_from_headers(&request_context.headers)?); - - validate_object_key(&key, if request_context.is_post_object { "POST" } else { "PUT" })?; - - let Some(body) = body else { return Err(s3_error!(IncompleteBody)) }; - - let mut size = resolved_size; - let mut transform_stage = PutObjectTransformStage::default(); - let mut plain_reduced_copy_stage = false; - let mut small_object_eager_stage = false; - let bytes_pool = get_concurrency_manager().bytes_pool(); - - let store = get_validated_store(&bucket).await?; - - 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(), - effective_kms_key_id.as_ref(), - sse_customer_algorithm.as_ref(), - sse_customer_key.as_ref(), - sse_customer_key_md5.as_ref(), - true, - )?; - - let encryption_enabled_for_put = effective_sse.is_some() - || effective_kms_key_id.is_some() - || sse_customer_algorithm.is_some() - || sse_customer_key.is_some() - || sse_customer_key_md5.is_some(); - - let body_plan = plan_put_object_body_with_transforms( - size, - &request_context.headers, - &key, - get_buffer_size_opt_in(size), - encryption_enabled_for_put, - ); - if body_plan.ingress.kind == PutObjectIngressKind::ReducedCopyCandidate { - rustfs_io_metrics::record_put_object_attempted_fast_path(size); - debug!( - encryption_enabled = encryption_enabled_for_put, - compressed = body_plan.should_compress(), - "Zero-copy write enabled for {} byte object (bucket={}, key={})", - size, - bucket, - key - ); - } else if let Some(reason) = resolve_put_transformed_fallback_reason( - body_plan.ingress.kind, - body_plan.should_compress(), - encryption_enabled_for_put, - ) { - rustfs_io_metrics::record_io_fallback(rustfs_io_metrics::IoStage::PutTransform, reason); - rustfs_io_metrics::record_put_fallback(size, reason); - } - - let mut metadata = metadata.unwrap_or_default(); - apply_put_request_metadata( - &mut metadata, - &request_context.headers, - &key, - cache_control, - content_disposition, - content_encoding, - content_language, - content_type, - expires, - website_redirect_location, - tagging, - storage_class.clone(), - )?; - - let mut opts: ObjectOptions = put_opts(&bucket, &key, version_id.clone(), &request_context.headers, metadata.clone()) - .await - .map_err(ApiError::from)?; - apply_put_request_object_lock_opts( - &bucket, - object_lock_legal_hold_status, - object_lock_mode, - object_lock_retain_until_date, - &mut opts, - ) - .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_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 - .map_err(ApiError::from)?; - match store.get_object_info(&bucket, &key, ¤t_opts).await { - Ok(existing_obj_info) => validate_existing_object_lock_for_write(&existing_obj_info)?, - Err(err) => { - if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) { - return Err(ApiError::from(err).into()); - } - } - } - - let actual_size = size; - let mut hash_values = PutObjectLegacyHashValues { - md5hex: if let Some(base64_md5) = content_md5 { - let md5 = base64_simd::STANDARD - .decode_to_vec(base64_md5.as_bytes()) - .map_err(|e| ApiError::from(StorageError::other(format!("Invalid content MD5: {e}"))))?; - Some(hex_simd::encode_to_string(&md5, hex_simd::AsciiCase::Lower)) - } else { - None - }, - sha256hex: get_content_sha256_with_query(&request_context.headers, request_context.uri_query.as_deref()), - }; - - let reader_stage_start = std::time::Instant::now(); - let stage = if can_use_small_put_eager - && should_use_small_put_eager_path(size, eager_max_bytes, body_plan.should_compress(), encryption_enabled_for_put) - { - small_object_eager_stage = true; - debug!( - "Plain PUT is using the eager small-object path (bucket={}, key={}, size={}, eager_max={})", - bucket, key, size, eager_max_bytes - ); - build_small_put_eager_hash_stage( - body, - size, - bytes_pool.clone(), - hash_values, - &request_context.headers, - request_context.trailing_headers.clone(), - ) - .await? - } else if body_plan.should_compress() { - transform_stage.mark_compression(); - let algorithm = CompressionAlgorithm::default(); - insert_str(&mut metadata, SUFFIX_COMPRESSION, algorithm.to_string()); - insert_str(&mut metadata, SUFFIX_ACTUAL_SIZE, size.to_string()); - - let ingress_source = build_put_object_ingress_source(body, body_plan); - let stage = build_put_object_plain_hash_stage( - ingress_source, - std::mem::take(&mut hash_values), - PutObjectLegacyHashStagePlan { - size, - actual_size: size, - apply_s3_checksum: true, - ignore_s3_checksum_value: false, - }, - &request_context.headers, - request_context.trailing_headers.clone(), - ) - .map_err(ApiError::from)?; - - if stage.ingress_kind == PutObjectIngressKind::ReducedCopyCandidate { - plain_reduced_copy_stage = true; - } - opts.want_checksum = stage.want_checksum; - insert_str(&mut opts.user_defined, SUFFIX_COMPRESSION, algorithm.to_string()); - insert_str(&mut opts.user_defined, SUFFIX_ACTUAL_SIZE, size.to_string()); - - let reader: Box = Box::new(CompressReader::new(stage.reader, algorithm)); - size = HashReader::SIZE_PRESERVE_LAYER; - hash_values.clear_for_transformed_body(); - build_put_object_legacy_hash_stage( - reader, - hash_values, - PutObjectLegacyHashStagePlan { - size, - actual_size, - apply_s3_checksum: size >= 0, - ignore_s3_checksum_value: false, - }, - &request_context.headers, - request_context.trailing_headers.clone(), - ) - .map_err(ApiError::from)? - } else { - let ingress_source = build_put_object_ingress_source(body, body_plan); - let stage = build_put_object_plain_hash_stage( - ingress_source, - hash_values, - PutObjectLegacyHashStagePlan { - size, - actual_size, - apply_s3_checksum: size >= 0, - ignore_s3_checksum_value: false, - }, - &request_context.headers, - request_context.trailing_headers.clone(), - ) - .map_err(ApiError::from)?; - - if stage.ingress_kind == PutObjectIngressKind::ReducedCopyCandidate { - plain_reduced_copy_stage = true; - debug!("Plain PUT is using the reduced-copy Reader hash path (bucket={}, key={})", bucket, key); - } - - 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, - put_path, - false, - ); - let mut reader = stage.reader; - if stage.want_checksum.is_some() { - opts.want_checksum = stage.want_checksum; - } - - let encryption_request = EncryptionRequest { - bucket: &bucket, - key: &key, - server_side_encryption: effective_sse.clone(), - ssekms_key_id: effective_kms_key_id.clone(), - sse_customer_algorithm: sse_customer_algorithm.clone(), - sse_customer_key, - sse_customer_key_md5: sse_customer_key_md5.clone(), - content_size: actual_size, - part_number: None, - part_key: None, - part_nonce: None, - }; - - if let Some(material) = sse_encryption(encryption_request).await? { - transform_stage.mark_encryption(); - effective_sse = Some(material.server_side_encryption.clone()); - effective_kms_key_id = material.kms_key_id.clone(); - - let encrypted_reader = material.wrap_reader(reader); - reader = HashReader::new(encrypted_reader, HashReader::SIZE_PRESERVE_LAYER, actual_size, None, None, false) - .map_err(ApiError::from)?; - - let encryption_metadata = material.metadata; - metadata.extend(encryption_metadata.clone()); - opts.user_defined.extend(encryption_metadata); - } - - let mut reader = PutObjReader::new(reader); - - let mt2 = metadata.clone(); - opts.user_defined.extend(metadata); - let capacity_scope_token = Uuid::new_v4(); - opts.capacity_scope_token = Some(capacity_scope_token); - - let repoptions = - get_must_replicate_options(&mt2, "".to_string(), ReplicationStatusType::Empty, ReplicationType::Object, opts.clone()); - - let dsc = must_replicate(&bucket, &key, repoptions).await; - - if dsc.replicate_any() { - insert_str(&mut opts.user_defined, SUFFIX_REPLICATION_TIMESTAMP, jiff::Zoned::now().to_string()); - insert_str( - &mut opts.user_defined, - SUFFIX_REPLICATION_STATUS, - dsc.pending_status().unwrap_or_default(), - ); - } - - let store_put_start = std::time::Instant::now(); - let obj_info = store - .put_object(&bucket, &key, &mut reader, &opts) - .await - .map_err(ApiError::from)?; - log_put_flow_phase( - &bucket, - &key, - "store_put_object", - store_put_start.elapsed(), - actual_size, - put_path, - transform_stage.encryption_applied(), - ); - - 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 BucketVersioningSys::prefix_enabled(&bucket, &key).await { - raw_version.clone() - } else { - None - }; - - let e_tag = obj_info.etag.clone().map(|etag| to_s3s_etag(&etag)); - - let repoptions = - get_must_replicate_options(&mt2, "".to_string(), ReplicationStatusType::Empty, ReplicationType::Object, opts); - - let dsc = must_replicate(&bucket, &key, repoptions).await; - let expiration = resolve_put_object_expiration(&bucket, &obj_info).await; - - if dsc.replicate_any() { - schedule_replication(obj_info.clone(), store.clone(), dsc, ReplicationType::Object).await; - } - - let mut checksums = PutObjectChecksums { - crc32: input.checksum_crc32, - crc32c: input.checksum_crc32c, - sha1: input.checksum_sha1, - sha256: input.checksum_sha256, - crc64nvme: input.checksum_crc64nvme, - }; - apply_trailing_checksums( - input.checksum_algorithm.as_ref().map(|a| a.as_str()), - &request_context.trailing_headers, - &mut checksums, - ); - checksums.merge_from_map(&reader.as_hash_reader().content_crc()); - if let Some(checksum_bytes) = resolved_checksum_bytes(&checksums) - && obj_info - .checksum - .as_ref() - .is_none_or(|stored| rustfs_rio::read_checksums(stored.as_ref(), 0).0.is_empty()) - { - let checksum_update_opts = ObjectOptions { - version_id: raw_version.clone(), - resolved_checksum: Some(checksum_bytes), - ..Default::default() - }; - let _ = store - .put_object_metadata(&bucket, &key, &checksum_update_opts) - .await - .map_err(ApiError::from)?; - } - - let output = PutObjectOutput { - e_tag, - server_side_encryption: effective_sse, - sse_customer_algorithm: sse_customer_algorithm.clone(), - sse_customer_key_md5: sse_customer_key_md5.clone(), - ssekms_key_id: effective_kms_key_id, - expiration, - checksum_crc32: checksums.crc32, - checksum_crc32c: checksums.crc32c, - checksum_sha1: checksums.sha1, - checksum_sha256: checksums.sha256, - checksum_crc64nvme: checksums.crc64nvme, - version_id: put_version, - ..Default::default() - }; - - record_capacity_write(Some(capacity_scope_token)).await; - - { - let duration_ms = start_time.elapsed().as_millis() as f64; - let fast_path_selected = plain_reduced_copy_stage || small_object_eager_stage; - rustfs_io_metrics::record_put_object(duration_ms, size, fast_path_selected); - let io_path = if fast_path_selected { - rustfs_io_metrics::IoPath::Fast - } else { - rustfs_io_metrics::IoPath::Legacy - }; - rustfs_io_metrics::record_io_path_selected("put", io_path); - rustfs_io_metrics::record_put_path_selected(actual_size, io_path); - let effective_copy_mode = transform_stage.effective_copy_mode(); - rustfs_io_metrics::record_io_copy_mode("put", effective_copy_mode, actual_size.max(0) as usize); - rustfs_io_metrics::record_put_copy_mode(actual_size, effective_copy_mode); - if let Some(transform_kind) = transform_stage.metric_kind() { - rustfs_io_metrics::record_put_transform_selected(transform_kind, io_path, actual_size.max(0) as usize); - } - } - - Ok((output, obj_info)) - } -} - -#[cfg(test)] -mod tests { - use super::*; - use bytes::Bytes; - use futures::{StreamExt, stream}; - use rustfs_io_core::BytesPool; - use serial_test::serial; - use std::sync::Arc; - use tokio::time::{Duration, timeout}; - - #[test] - fn small_put_eager_path_only_targets_plain_small_objects() { - assert!(should_use_small_put_eager_path( - 64 * 1024, - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, - false, - false - )); - assert!(should_use_small_put_eager_path( - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, - false, - false - )); - assert!(!should_use_small_put_eager_path( - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES + 1, - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, - false, - false - )); - assert!(!should_use_small_put_eager_path( - 64 * 1024, - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, - true, - false - )); - assert!(!should_use_small_put_eager_path( - 64 * 1024, - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, - false, - true - )); - assert!(!should_use_small_put_eager_path(0, DEFAULT_SMALL_PUT_EAGER_MAX_BYTES, false, false)); - } - - #[test] - fn clamp_small_put_eager_max_bytes_caps_inline_budget() { - assert_eq!( - clamp_small_put_eager_max_bytes(Some(rustfs_object_io::put::PUT_REDUCED_COPY_MIN_SIZE_BYTES as usize * 2)), - DEFAULT_SMALL_PUT_EAGER_MAX_BYTES - ); - assert_eq!(clamp_small_put_eager_max_bytes(Some(128 * 1024)), 128 * 1024); - assert_eq!(clamp_small_put_eager_max_bytes(None), DEFAULT_SMALL_PUT_EAGER_MAX_BYTES); - } - - #[test] - #[serial] - fn resolved_small_put_eager_max_bytes_honors_disable_env() { - temp_env::with_var(ENV_RUSTFS_PUT_FORCE_DISABLE_SMALL_EAGER, Some("true"), || { - assert_eq!(resolved_small_put_eager_max_bytes(256 * 1024), 0); - }); - } - - #[test] - #[serial] - fn resolved_small_put_eager_max_bytes_narrows_default_budget() { - temp_env::with_var(ENV_RUSTFS_PUT_SMALL_EAGER_MAX_BYTES, Some("4096"), || { - assert_eq!(resolved_small_put_eager_max_bytes(256 * 1024), 4096); - }); - } - - #[test] - #[serial] - fn resolved_small_put_eager_max_bytes_ignores_invalid_override() { - temp_env::with_var(ENV_RUSTFS_PUT_SMALL_EAGER_MAX_BYTES, Some("invalid"), || { - assert_eq!(resolved_small_put_eager_max_bytes(256 * 1024), 256 * 1024); - }); - } - - #[tokio::test] - async fn read_small_put_body_eager_requires_exact_content_length() { - let body = stream::iter(vec![ - Ok::(Bytes::from_static(b"abc")), - Ok::(Bytes::from_static(b"def")), - ]); - - let pool = Arc::new(BytesPool::new_tiered()); - let data = read_small_put_body_eager(body, 6, pool) - .await - .expect("eager read should succeed"); - assert_eq!(data.as_ref(), b"abcdef"); - } - - #[tokio::test] - async fn read_small_put_body_eager_rejects_length_mismatch() { - let body = stream::iter(vec![Ok::(Bytes::from_static(b"abc"))]); - let pool = Arc::new(BytesPool::new_tiered()); - - let err = read_small_put_body_eager(body, 4, pool) - .await - .expect_err("short eager read should fail"); - assert_eq!(err.code(), &S3ErrorCode::IncompleteBody); - } - - #[tokio::test] - async fn read_small_put_body_eager_rejects_overlong_body() { - let body = stream::iter(vec![ - Ok::(Bytes::from_static(b"abc")), - Ok::(Bytes::from_static(b"def")), - ]); - let pool = Arc::new(BytesPool::new_tiered()); - - let err = read_small_put_body_eager(body, 5, pool) - .await - .expect_err("overlong eager read should fail"); - assert_eq!(err.code(), &S3ErrorCode::IncompleteBody); - } - - #[tokio::test] - async fn read_small_put_body_eager_returns_buffer_to_pool_after_drop() { - let body = stream::iter(vec![Ok::(Bytes::from_static(b"abc"))]); - let pool = Arc::new(BytesPool::new_tiered()); - - let data = read_small_put_body_eager(body, 3, pool.clone()) - .await - .expect("pooled eager read should succeed"); - assert_eq!(pool.available_buffers(), 0); - - drop(data); - assert_eq!(pool.available_buffers(), 1); - } - - #[tokio::test] - async fn read_small_put_body_eager_returns_after_expected_bytes_without_waiting_for_eof() { - let body = stream::once(async { Ok::(Bytes::from_static(b"abc")) }).chain(stream::pending()); - let pool = Arc::new(BytesPool::new_tiered()); - - let data = timeout(Duration::from_millis(50), read_small_put_body_eager(body, 3, pool)) - .await - .expect("eager read should not wait for stream termination") - .expect("eager read should succeed once content-length bytes are read"); - - assert_eq!(data.as_ref(), b"abc"); - } -}