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