refactor(storage): centralize S3 response/error helpers (#1818)

This commit is contained in:
安正超
2026-02-14 21:20:57 +08:00
committed by GitHub
parent 546485a8ee
commit d3ff6ff36a
3 changed files with 200 additions and 125 deletions
+121 -124
View File
@@ -26,6 +26,9 @@ use crate::storage::readers::InMemoryAsyncReader;
use crate::storage::s3_api::bucket::{build_list_objects_output, build_list_objects_v2_output};
use crate::storage::s3_api::common::rustfs_owner;
use crate::storage::s3_api::multipart::{build_list_multipart_uploads_output, build_list_parts_output, parse_list_parts_params};
use crate::storage::s3_api::response::{
access_denied_error, map_abort_multipart_upload_error, not_initialized_error, s3_response,
};
use crate::storage::sse::{
DecryptionRequest, EncryptionRequest, PrepareEncryptionRequest, check_encryption_metadata, sse_decryption, sse_encryption,
sse_prepare_encryption, strip_managed_encryption_metadata,
@@ -267,7 +270,7 @@ impl FS {
})?;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let prefix = req
@@ -352,7 +355,7 @@ impl FS {
bucket_name: bucket.clone(),
object: _obj_info.clone(),
req_params: extract_params_header(&req.headers),
resp_elements: extract_resp_elements(&S3Response::new(output.clone())),
resp_elements: extract_resp_elements(&s3_response(output.clone())),
version_id: version_id.clone(),
host: get_request_host(&req.headers),
port: get_request_port(&req.headers),
@@ -428,7 +431,7 @@ impl FS {
checksum_crc64nvme,
..Default::default()
};
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -465,7 +468,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let opts = &ObjectOptions::default();
@@ -480,14 +483,8 @@ impl S3 for FS {
.abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts)
.await
{
Ok(_) => Ok(S3Response::new(AbortMultipartUploadOutput { ..Default::default() })),
Err(err) => {
// Convert MalformedUploadID to NoSuchUpload for S3 API compatibility
if matches!(err, StorageError::MalformedUploadID(_)) {
return Err(S3Error::new(S3ErrorCode::NoSuchUpload));
}
Err(ApiError::from(err).into())
}
Ok(_) => Ok(s3_response(AbortMultipartUploadOutput { ..Default::default() })),
Err(err) => Err(map_abort_multipart_upload_error(err)),
}
}
@@ -511,7 +508,7 @@ impl S3 for FS {
if if_match.is_some() || if_none_match.is_some() {
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
match store.get_object_info(&bucket, &key, &ObjectOptions::default()).await {
@@ -582,7 +579,7 @@ impl S3 for FS {
// TODO: check object lock
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// TDD: Get multipart info to extract encryption configuration before completing
@@ -764,9 +761,9 @@ impl S3 for FS {
helper = helper.version_id(version_id.clone());
}
let helper_result = Ok(S3Response::new(helper_output));
let helper_result = Ok(s3_response(helper_output));
let _ = helper.complete(&helper_result);
Ok(S3Response::new(output))
Ok(s3_response(output))
}
/// Copy an object from one location to another
@@ -839,7 +836,7 @@ impl S3 for FS {
}
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let h = HeaderMap::new();
@@ -1075,7 +1072,7 @@ impl S3 for FS {
let version_id = req.input.version_id.clone().unwrap_or_default();
helper = helper.object(object_info).version_id(version_id);
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -1094,7 +1091,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
counter!("rustfs_create_bucket_total").increment(1);
@@ -1113,7 +1110,7 @@ impl S3 for FS {
let output = CreateBucketOutput::default();
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -1149,7 +1146,7 @@ impl S3 for FS {
// debug!("create_multipart_upload meta {:?}", &metadata);
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let mut metadata = extract_metadata(&req.headers);
@@ -1224,7 +1221,7 @@ impl S3 for FS {
..Default::default()
};
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -1236,7 +1233,7 @@ impl S3 for FS {
let input = req.input.clone();
// TODO: DeleteBucketInput doesn't have force parameter?
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// get value from header, support mc style
@@ -1268,7 +1265,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
let result = Ok(S3Response::new(DeleteBucketOutput {}));
let result = Ok(s3_response(DeleteBucketOutput {}));
let _ = helper.complete(&result);
result
}
@@ -1278,7 +1275,7 @@ impl S3 for FS {
let DeleteBucketCorsInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -1290,7 +1287,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(DeleteBucketCorsOutput {}))
Ok(s3_response(DeleteBucketCorsOutput {}))
}
async fn delete_bucket_encryption(
@@ -1300,7 +1297,7 @@ impl S3 for FS {
let DeleteBucketEncryptionInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -1311,7 +1308,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(DeleteBucketEncryptionOutput::default()))
Ok(s3_response(DeleteBucketEncryptionOutput::default()))
}
#[instrument(level = "debug", skip(self))]
@@ -1322,7 +1319,7 @@ impl S3 for FS {
let DeleteBucketLifecycleInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -1334,7 +1331,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(DeleteBucketLifecycleOutput::default()))
Ok(s3_response(DeleteBucketLifecycleOutput::default()))
}
async fn delete_bucket_policy(
@@ -1344,7 +1341,7 @@ impl S3 for FS {
let DeleteBucketPolicyInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -1356,7 +1353,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(DeleteBucketPolicyOutput {}))
Ok(s3_response(DeleteBucketPolicyOutput {}))
}
async fn delete_bucket_replication(
@@ -1366,7 +1363,7 @@ impl S3 for FS {
let DeleteBucketReplicationInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -1380,7 +1377,7 @@ impl S3 for FS {
// TODO: remove targets
error!("delete bucket");
Ok(S3Response::new(DeleteBucketReplicationOutput::default()))
Ok(s3_response(DeleteBucketReplicationOutput::default()))
}
#[instrument(level = "debug", skip(self))]
@@ -1394,7 +1391,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(DeleteBucketTaggingOutput {}))
Ok(s3_response(DeleteBucketTaggingOutput {}))
}
/// Delete an object
@@ -1445,7 +1442,7 @@ impl S3 for FS {
}
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// Check Object Lock retention before deletion
@@ -1552,7 +1549,7 @@ impl S3 for FS {
.object(obj_info)
.version_id(version_id.map(|v| v.to_string()).unwrap_or_default());
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -1573,7 +1570,7 @@ impl S3 for FS {
let Some(store) = new_object_layer_fn() else {
error!("Store not initialized");
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// Support versioned objects
@@ -1608,7 +1605,7 @@ impl S3 for FS {
let version_id_resp = version_id.clone().unwrap_or_default();
helper = helper.version_id(version_id_resp);
let result = Ok(S3Response::new(DeleteObjectTaggingOutput { version_id }));
let result = Ok(s3_response(DeleteObjectTaggingOutput { version_id }));
let _ = helper.complete(&result);
let duration = start_time.elapsed();
histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "delete").record(duration.as_secs_f64());
@@ -1646,7 +1643,7 @@ impl S3 for FS {
.await;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let version_cfg = BucketVersioningSys::get(&bucket).await.unwrap_or_default();
@@ -1928,7 +1925,7 @@ impl S3 for FS {
)
.version_id(dobj.version_id.map(|v| v.to_string()).unwrap_or_default())
.req_params(extract_params_header(&req_headers))
.resp_elements(extract_resp_elements(&S3Response::new(DeleteObjectsOutput::default())))
.resp_elements(extract_resp_elements(&s3_response(DeleteObjectsOutput::default())))
.host(get_request_host(&req_headers))
.user_agent(get_request_user_agent(&req_headers))
.build();
@@ -1938,7 +1935,7 @@ impl S3 for FS {
}
});
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -1947,7 +1944,7 @@ impl S3 for FS {
let GetBucketAclInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -1966,7 +1963,7 @@ impl S3 for FS {
permission: Some(Permission::from_static(Permission::FULL_CONTROL)),
}];
Ok(S3Response::new(GetBucketAclOutput {
Ok(s3_response(GetBucketAclOutput {
grants: Some(grants),
owner: Some(RUSTFS_OWNER.to_owned()),
}))
@@ -1997,7 +1994,7 @@ impl S3 for FS {
}
};
Ok(S3Response::new(GetBucketCorsOutput {
Ok(s3_response(GetBucketCorsOutput {
cors_rules: Some(cors_configuration.cors_rules),
}))
}
@@ -2009,7 +2006,7 @@ impl S3 for FS {
let GetBucketEncryptionInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2028,7 +2025,7 @@ impl S3 for FS {
}
};
Ok(S3Response::new(GetBucketEncryptionOutput {
Ok(s3_response(GetBucketEncryptionOutput {
server_side_encryption_configuration,
}))
}
@@ -2041,7 +2038,7 @@ impl S3 for FS {
let GetBucketLifecycleConfigurationInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2058,7 +2055,7 @@ impl S3 for FS {
}
};
Ok(S3Response::new(GetBucketLifecycleConfigurationOutput {
Ok(s3_response(GetBucketLifecycleConfigurationOutput {
rules: Some(rules),
..Default::default()
}))
@@ -2071,7 +2068,7 @@ impl S3 for FS {
let input = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2080,13 +2077,13 @@ impl S3 for FS {
.map_err(ApiError::from)?;
if let Some(region) = rustfs_ecstore::global::get_global_region() {
return Ok(S3Response::new(GetBucketLocationOutput {
return Ok(s3_response(GetBucketLocationOutput {
location_constraint: Some(BucketLocationConstraint::from(region)),
}));
}
let output = GetBucketLocationOutput::default();
Ok(S3Response::new(output))
Ok(s3_response(output))
}
async fn get_bucket_notification_configuration(
@@ -2096,7 +2093,7 @@ impl S3 for FS {
let GetBucketNotificationConfigurationInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2118,14 +2115,14 @@ impl S3 for FS {
topic_configurations,
}) = has_notification_config
{
Ok(S3Response::new(GetBucketNotificationConfigurationOutput {
Ok(s3_response(GetBucketNotificationConfigurationOutput {
event_bridge_configuration,
lambda_function_configurations,
queue_configurations,
topic_configurations,
}))
} else {
Ok(S3Response::new(GetBucketNotificationConfigurationOutput::default()))
Ok(s3_response(GetBucketNotificationConfigurationOutput::default()))
}
}
@@ -2133,7 +2130,7 @@ impl S3 for FS {
let GetBucketPolicyInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2153,7 +2150,7 @@ impl S3 for FS {
}
};
Ok(S3Response::new(GetBucketPolicyOutput {
Ok(s3_response(GetBucketPolicyOutput {
policy: Some(policy_str),
}))
}
@@ -2165,7 +2162,7 @@ impl S3 for FS {
let GetBucketPolicyStatusInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2206,7 +2203,7 @@ impl S3 for FS {
}),
};
Ok(S3Response::new(output))
Ok(s3_response(output))
}
async fn get_bucket_replication(
@@ -2216,7 +2213,7 @@ impl S3 for FS {
let GetBucketReplicationInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2245,12 +2242,12 @@ impl S3 for FS {
));
}
// Ok(S3Response::new(GetBucketReplicationOutput {
// Ok(s3_response(GetBucketReplicationOutput {
// replication_configuration: rcfg,
// }))
if rcfg.is_some() {
Ok(S3Response::new(GetBucketReplicationOutput {
Ok(s3_response(GetBucketReplicationOutput {
replication_configuration: rcfg,
}))
} else {
@@ -2258,7 +2255,7 @@ impl S3 for FS {
role: "".to_string(),
rules: vec![],
};
Ok(S3Response::new(GetBucketReplicationOutput {
Ok(s3_response(GetBucketReplicationOutput {
replication_configuration: Some(rep),
}))
}
@@ -2286,7 +2283,7 @@ impl S3 for FS {
}
};
Ok(S3Response::new(GetBucketTaggingOutput { tag_set }))
Ok(s3_response(GetBucketTaggingOutput { tag_set }))
}
#[instrument(level = "debug", skip(self))]
@@ -2296,7 +2293,7 @@ impl S3 for FS {
) -> S3Result<S3Response<GetBucketVersioningOutput>> {
let GetBucketVersioningInput { bucket, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -2306,7 +2303,7 @@ impl S3 for FS {
let VersioningConfiguration { status, .. } = BucketVersioningSys::get(&bucket).await.map_err(ApiError::from)?;
Ok(S3Response::new(GetBucketVersioningOutput {
Ok(s3_response(GetBucketVersioningOutput {
status,
..Default::default()
}))
@@ -2444,7 +2441,7 @@ impl S3 for FS {
// S3 bucket notifications (s3:GetObject events) are triggered.
// This ensures event-driven workflows (Lambda, SNS) work correctly
// for both cache hits and misses.
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
return result;
}
@@ -2891,7 +2888,7 @@ impl S3 for FS {
let GetObjectAclInput { bucket, key, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
if let Err(e) = store.get_object_info(&bucket, &key, &ObjectOptions::default()).await {
@@ -2909,7 +2906,7 @@ impl S3 for FS {
permission: Some(Permission::from_static(Permission::FULL_CONTROL)),
}];
Ok(S3Response::new(GetObjectAclOutput {
Ok(s3_response(GetObjectAclOutput {
grants: Some(grants),
owner: Some(RUSTFS_OWNER.to_owned()),
..Default::default()
@@ -2924,7 +2921,7 @@ impl S3 for FS {
let GetObjectAttributesInput { bucket, key, .. } = req.input.clone();
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
if let Err(e) = store
@@ -2949,7 +2946,7 @@ impl S3 for FS {
})
.version_id(version_id);
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -2964,7 +2961,7 @@ impl S3 for FS {
} = req.input.clone();
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let _ = store
@@ -3004,7 +3001,7 @@ impl S3 for FS {
let version_id = req.input.version_id.clone().unwrap_or_else(|| Uuid::new_v4().to_string());
helper = helper.object(object_info).version_id(version_id);
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -3035,7 +3032,7 @@ impl S3 for FS {
// warn!("object_lock_configuration {:?}", &object_lock_configuration);
Ok(S3Response::new(GetObjectLockConfigurationOutput {
Ok(s3_response(GetObjectLockConfigurationOutput {
object_lock_configuration,
}))
}
@@ -3050,7 +3047,7 @@ impl S3 for FS {
} = req.input.clone();
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// check object lock
@@ -3082,7 +3079,7 @@ impl S3 for FS {
let version_id = req.input.version_id.clone().unwrap_or_default();
helper = helper.object(object_info).version_id(version_id);
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -3096,7 +3093,7 @@ impl S3 for FS {
let Some(store) = new_object_layer_fn() else {
error!("Store not initialized");
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// Support versioned objects
@@ -3122,7 +3119,7 @@ impl S3 for FS {
counter!("rustfs.get_object_tagging.success").increment(1);
let duration = start_time.elapsed();
histogram!("rustfs.object_tagging.operation.duration.seconds", "operation" => "put").record(duration.as_secs_f64());
Ok(S3Response::new(GetObjectTaggingOutput {
Ok(s3_response(GetObjectTaggingOutput {
tag_set,
version_id: req.input.version_id.clone(),
}))
@@ -3141,7 +3138,7 @@ impl S3 for FS {
let input = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -3150,7 +3147,7 @@ impl S3 for FS {
.map_err(ApiError::from)?;
// mc cp step 2 GetBucketInfo
Ok(S3Response::new(HeadBucketOutput::default()))
Ok(s3_response(HeadBucketOutput::default()))
}
#[instrument(level = "debug", skip(self, req))]
@@ -3197,7 +3194,7 @@ impl S3 for FS {
.map_err(ApiError::from)?;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// Modification Points: Explicitly handles get_object_info errors, distinguishing between object absence and other errors
let info = match store.get_object_info(&bucket, &key, &opts).await {
@@ -3492,13 +3489,13 @@ impl S3 for FS {
// mc ls
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let mut req = req;
if req.credentials.as_ref().is_none_or(|cred| cred.access_key.is_empty()) {
return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied"));
return Err(access_denied_error());
}
let bucket_infos = if let Err(e) = authorize_request(&mut req, Action::S3Action(S3Action::ListAllMyBucketsAction)).await {
@@ -3530,7 +3527,7 @@ impl S3 for FS {
.await;
if list_bucket_infos.is_empty() {
return Err(S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied"));
return Err(access_denied_error());
}
list_bucket_infos
} else {
@@ -3551,7 +3548,7 @@ impl S3 for FS {
owner: Some(RUSTFS_OWNER.to_owned()),
..Default::default()
};
Ok(S3Response::new(output))
Ok(s3_response(output))
}
async fn list_multipart_uploads(
@@ -3569,7 +3566,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let prefix = prefix.unwrap_or_default();
@@ -3589,7 +3586,7 @@ impl S3 for FS {
let output = build_list_multipart_uploads_output(bucket, prefix, result);
Ok(S3Response::new(output))
Ok(s3_response(output))
}
async fn list_object_versions(
@@ -3678,7 +3675,7 @@ impl S3 for FS {
..Default::default()
};
Ok(S3Response::new(output))
Ok(s3_response(output))
}
#[instrument(level = "debug", skip(self, req))]
@@ -3783,7 +3780,7 @@ impl S3 for FS {
response_start_after,
);
Ok(S3Response::new(output))
Ok(s3_response(output))
}
#[instrument(level = "debug", skip(self, req))]
@@ -3798,7 +3795,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let (part_number_marker, max_parts) = parse_list_parts_params(part_number_marker, max_parts)?;
@@ -3809,7 +3806,7 @@ impl S3 for FS {
.map_err(ApiError::from)?;
let output = build_list_parts_output(res);
Ok(S3Response::new(output))
Ok(s3_response(output))
}
async fn put_bucket_acl(&self, req: S3Request<PutBucketAclInput>) -> S3Result<S3Response<PutBucketAclOutput>> {
@@ -3823,7 +3820,7 @@ impl S3 for FS {
// TODO:checkRequestAuthType
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -3852,7 +3849,7 @@ impl S3 for FS {
return Err(s3_error!(NotImplemented));
}
}
Ok(S3Response::new(PutBucketAclOutput::default()))
Ok(s3_response(PutBucketAclOutput::default()))
}
#[instrument(level = "debug", skip(self))]
@@ -3864,7 +3861,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -3878,7 +3875,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(PutBucketCorsOutput::default()))
Ok(s3_response(PutBucketCorsOutput::default()))
}
async fn put_bucket_encryption(
@@ -3894,7 +3891,7 @@ impl S3 for FS {
info!("sse_config {:?}", &server_side_encryption_configuration);
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -3908,7 +3905,7 @@ impl S3 for FS {
metadata_sys::update(&bucket, BUCKET_SSECONFIG, data)
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(PutBucketEncryptionOutput::default()))
Ok(s3_response(PutBucketEncryptionOutput::default()))
}
#[instrument(level = "debug", skip(self))]
@@ -3942,7 +3939,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(PutBucketLifecycleConfigurationOutput::default()))
Ok(s3_response(PutBucketLifecycleConfigurationOutput::default()))
}
async fn put_bucket_notification_configuration(
@@ -3956,7 +3953,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// Verify that the bucket exists
@@ -4014,14 +4011,14 @@ impl S3 for FS {
.await
.map_err(|e| s3_error!(InternalError, "Failed to add rules: {e}"))?;
Ok(S3Response::new(PutBucketNotificationConfigurationOutput {}))
Ok(s3_response(PutBucketNotificationConfigurationOutput {}))
}
async fn put_bucket_policy(&self, req: S3Request<PutBucketPolicyInput>) -> S3Result<S3Response<PutBucketPolicyOutput>> {
let PutBucketPolicyInput { bucket, policy, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -4047,7 +4044,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(PutBucketPolicyOutput {}))
Ok(s3_response(PutBucketPolicyOutput {}))
}
async fn put_bucket_replication(
@@ -4062,7 +4059,7 @@ impl S3 for FS {
warn!("put bucket replication");
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -4077,7 +4074,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(PutBucketReplicationOutput::default()))
Ok(s3_response(PutBucketReplicationOutput::default()))
}
#[instrument(level = "debug", skip(self))]
@@ -4085,7 +4082,7 @@ impl S3 for FS {
let PutBucketTaggingInput { bucket, tagging, .. } = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -4099,7 +4096,7 @@ impl S3 for FS {
.await
.map_err(ApiError::from)?;
Ok(S3Response::new(Default::default()))
Ok(s3_response(Default::default()))
}
#[instrument(level = "debug", skip(self))]
@@ -4126,7 +4123,7 @@ impl S3 for FS {
// TODO: globalSiteReplicationSys.BucketMetaHook
Ok(S3Response::new(PutBucketVersioningOutput {}))
Ok(s3_response(PutBucketVersioningOutput {}))
}
#[instrument(level = "debug", skip(self, req))]
@@ -4144,7 +4141,7 @@ impl S3 for FS {
} = req.input;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
if let Err(e) = store.get_object_info(&bucket, &key, &ObjectOptions::default()).await {
@@ -4172,7 +4169,7 @@ impl S3 for FS {
return Err(s3_error!(NotImplemented));
}
}
Ok(S3Response::new(PutObjectAclOutput::default()))
Ok(s3_response(PutObjectAclOutput::default()))
}
async fn put_object_legal_hold(
@@ -4189,7 +4186,7 @@ impl S3 for FS {
} = req.input.clone();
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let _ = store
@@ -4223,7 +4220,7 @@ impl S3 for FS {
let version_id = req.input.version_id.clone().unwrap_or_default();
helper = helper.object(info).version_id(version_id);
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -4242,7 +4239,7 @@ impl S3 for FS {
let Some(input_cfg) = object_lock_configuration else { return Err(s3_error!(InvalidArgument)) };
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
store
@@ -4287,7 +4284,7 @@ impl S3 for FS {
.map_err(ApiError::from)?;
}
Ok(S3Response::new(PutObjectLockConfigurationOutput::default()))
Ok(s3_response(PutObjectLockConfigurationOutput::default()))
}
async fn put_object_retention(
@@ -4304,7 +4301,7 @@ impl S3 for FS {
} = req.input.clone();
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// check object lock
@@ -4365,7 +4362,7 @@ impl S3 for FS {
let version_id = req.input.version_id.clone().unwrap_or_else(|| Uuid::new_v4().to_string());
helper = helper.object(object_info).version_id(version_id);
let result = Ok(S3Response::new(output));
let result = Ok(s3_response(output));
let _ = helper.complete(&result);
result
}
@@ -4391,7 +4388,7 @@ impl S3 for FS {
}
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let mut tag_keys = std::collections::HashSet::with_capacity(tagging.tag_set.len());
@@ -4458,7 +4455,7 @@ impl S3 for FS {
let version_id_resp = req.input.version_id.clone().unwrap_or_default();
helper = helper.version_id(version_id_resp);
let result = Ok(S3Response::new(PutObjectTaggingOutput {
let result = Ok(s3_response(PutObjectTaggingOutput {
version_id: req.input.version_id.clone(),
}));
let _ = helper.complete(&result);
@@ -4481,7 +4478,7 @@ impl S3 for FS {
})?;
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let version_id_str = version_id.clone().unwrap_or_default();
@@ -4583,7 +4580,7 @@ impl S3 for FS {
request_charged: Some(RequestCharged::from_static(RequestCharged::REQUESTER)),
restore_output_path: None,
};
return Ok(S3Response::new(output));
return Ok(s3_response(output));
}
}
@@ -4706,7 +4703,7 @@ impl S3 for FS {
drop(tx);
});
Ok(S3Response::new(SelectObjectContentOutput {
Ok(s3_response(SelectObjectContentOutput {
payload: Some(SelectObjectContentEventStream::new(stream)),
}))
}
@@ -4764,7 +4761,7 @@ impl S3 for FS {
// Get multipart info early to check if managed encryption will be applied
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
let opts = ObjectOptions::default();
@@ -4954,7 +4951,7 @@ impl S3 for FS {
..Default::default()
};
Ok(S3Response::new(output))
Ok(s3_response(output))
}
#[instrument(level = "debug", skip(self, req))]
@@ -4996,7 +4993,7 @@ impl S3 for FS {
// Note: In a real implementation, you would properly validate access
// For now, we'll skip the detailed authorization check
let Some(store) = new_object_layer_fn() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
return Err(not_initialized_error());
};
// Check if multipart upload exists and get its info
@@ -5217,6 +5214,6 @@ impl S3 for FS {
..Default::default()
};
Ok(S3Response::new(output))
Ok(s3_response(output))
}
}
+1 -1
View File
@@ -29,7 +29,7 @@ pub(crate) mod multipart;
/// Object-specific extraction steps can be added here incrementally.
pub(crate) mod object {}
pub(crate) mod replication {}
pub(crate) mod response {}
pub(crate) mod response;
pub(crate) mod restore {}
pub(crate) mod select {}
pub(crate) mod validation {}
+78
View File
@@ -0,0 +1,78 @@
// 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 crate::error::ApiError;
use rustfs_ecstore::error::StorageError;
use s3s::{S3Error, S3ErrorCode, S3Response};
pub(crate) fn s3_response<T>(output: T) -> S3Response<T> {
S3Response::new(output)
}
pub(crate) fn not_initialized_error() -> S3Error {
S3Error::with_message(S3ErrorCode::InternalError, "Not init")
}
pub(crate) fn access_denied_error() -> S3Error {
S3Error::with_message(S3ErrorCode::AccessDenied, "Access Denied")
}
pub(crate) fn map_abort_multipart_upload_error(err: StorageError) -> S3Error {
// For abort multipart upload, malformed upload IDs should be hidden as NoSuchUpload
// to match S3 API compatibility expectations.
if matches!(err, StorageError::MalformedUploadID(_)) {
return S3Error::new(S3ErrorCode::NoSuchUpload);
}
ApiError::from(err).into()
}
#[cfg(test)]
mod tests {
use super::{access_denied_error, map_abort_multipart_upload_error, not_initialized_error, s3_response};
use rustfs_ecstore::error::StorageError;
use s3s::{S3ErrorCode, S3Response};
#[test]
fn test_s3_response_wraps_output() {
let response: S3Response<i32> = s3_response(7);
assert_eq!(response.output, 7);
}
#[test]
fn test_not_initialized_error_shape() {
let err = not_initialized_error();
assert_eq!(*err.code(), S3ErrorCode::InternalError);
assert_eq!(err.message(), Some("Not init"));
}
#[test]
fn test_access_denied_error_shape() {
let err = access_denied_error();
assert_eq!(*err.code(), S3ErrorCode::AccessDenied);
assert_eq!(err.message(), Some("Access Denied"));
}
#[test]
fn test_map_abort_multipart_upload_error_for_malformed_id() {
let err = map_abort_multipart_upload_error(StorageError::MalformedUploadID("bad-id".to_string()));
assert_eq!(*err.code(), S3ErrorCode::NoSuchUpload);
}
#[test]
fn test_map_abort_multipart_upload_error_for_unexpected_error() {
let err = map_abort_multipart_upload_error(StorageError::Unexpected);
assert_eq!(*err.code(), S3ErrorCode::InternalError);
}
}