cleanup: remove orphaned chunk fast-path submodule files (#2545)

Co-authored-by: cxymds <Cxymds@qq.com>
This commit is contained in:
安正超
2026-04-15 11:50:32 +08:00
committed by GitHub
parent 1eef75b0ca
commit df109680ae
4 changed files with 0 additions and 1999 deletions
@@ -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<R>(
final_stream: R,
bucket: &str,
key: &str,
response_content_length: i64,
optimal_buffer_size: usize,
planning_inputs: ObjectIoGetObjectBodyPlanningInputs,
) -> S3Result<Option<StreamingBlob>>
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<GetObjectFlowResult> {
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<dyn Reader>;
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);
}
}
@@ -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<GetObjectIoPlanning<'a>> {
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<GetObjectReadSetup> {
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<dyn Reader>,
),
};
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<GetObjectPreparedRead<'a>> {
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 })
}
@@ -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<dyn NotifyInterface>,
resolved_size: i64,
) -> S3Result<PutObjectOutput> {
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::<request_context::RequestContext>().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<T>(input: T, method: Method) -> S3Request<T> {
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);
}
}
@@ -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<Bytes> {
[
(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<usize>) -> 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<std::io::Result<()>> {
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<S, B, E>(body: S, size: i64, pool: Arc<BytesPool>) -> S3Result<PooledBuffer>
where
S: Stream<Item = Result<B, E>>,
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<S, B, E>(
body: S,
size: i64,
pool: Arc<BytesPool>,
hash_values: PutObjectLegacyHashValues,
headers: &HeaderMap,
trailing_headers: Option<s3s::TrailingHeaders>,
) -> S3Result<rustfs_object_io::put::PutObjectHashStage>
where
S: Stream<Item = Result<B, E>>,
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, &current_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<dyn Reader> = 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, std::io::Error>(Bytes::from_static(b"abc")),
Ok::<Bytes, std::io::Error>(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, std::io::Error>(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, std::io::Error>(Bytes::from_static(b"abc")),
Ok::<Bytes, std::io::Error>(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, std::io::Error>(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, std::io::Error>(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");
}
}