refactor(storage): inline extract put-object branch (#2482)

This commit is contained in:
安正超
2026-04-11 11:00:11 +08:00
committed by GitHub
parent d70ca1990e
commit b49570d87a
2 changed files with 37 additions and 26 deletions
+17 -26
View File
@@ -719,7 +719,23 @@ impl DefaultObjectUsecase {
return Err(s3_error!(InvalidStorageClass));
}
if is_put_object_extract_requested(&request_context.headers) {
return self.execute_put_object_extract(req, request_context).await;
let helper = OperationHelper::new(&req, event_name, S3Operation::PutObject).suppress_event();
if is_sse_kms_requested(&req.input, &request_context.headers) {
return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for extract uploads"));
}
let resolved_size = resolve_put_body_size(req.input.content_length, &request_context.headers)?;
self.check_bucket_quota(&req.input.bucket, quota_operation, resolved_size as u64)
.await?;
let notify = self
.context
.as_ref()
.map(|context| context.notify())
.unwrap_or_else(default_notify_interface);
let input = req.input;
let output = DefaultObjectUsecase::run_put_object_extract_flow(input, request_context, notify, resolved_size).await?;
let result = Ok(S3Response::new(output));
let _ = helper.complete(&result);
return result;
}
let helper = OperationHelper::new(&req, event_name, S3Operation::PutObject);
@@ -3025,31 +3041,6 @@ impl DefaultObjectUsecase {
payload: Some(SelectObjectContentEventStream::new(stream)),
}))
}
#[instrument(level = "debug", skip(self, req, request_context))]
async fn execute_put_object_extract(
&self,
req: S3Request<PutObjectInput>,
request_context: PutObjectRequestContext,
) -> S3Result<S3Response<PutObjectOutput>> {
let helper = OperationHelper::new(&req, EventName::ObjectCreatedPut, S3Operation::PutObject).suppress_event();
if is_sse_kms_requested(&req.input, &request_context.headers) {
return Err(s3_error!(NotImplemented, "SSE-KMS is not supported for extract uploads"));
}
let resolved_size = resolve_put_body_size(req.input.content_length, &request_context.headers)?;
self.check_bucket_quota(&req.input.bucket, QuotaOperation::PutObject, resolved_size as u64)
.await?;
let notify = self
.context
.as_ref()
.map(|context| context.notify())
.unwrap_or_else(default_notify_interface);
let input = req.input;
let output = DefaultObjectUsecase::run_put_object_extract_flow(input, request_context, notify, resolved_size).await?;
let result = Ok(S3Response::new(output));
let _ = helper.complete(&result);
result
}
}
fn object_attributes_requested(object_attributes: &[ObjectAttributes], name: &'static str) -> bool {
@@ -491,4 +491,24 @@ mod tests {
let err = usecase.execute_put_object(&fs, req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::NotImplemented);
}
#[tokio::test]
async fn execute_put_object_rejects_extract_sse_kms_key_id_header() {
let input = PutObjectInput::builder()
.bucket("test-bucket".to_string())
.key("archive.tar".to_string())
.build()
.unwrap();
let mut req = build_request(input, Method::PUT);
req.headers.insert(AMZ_SNOWBALL_EXTRACT, HeaderValue::from_static("true"));
req.headers
.insert(AMZ_SERVER_SIDE_ENCRYPTION_KMS_ID, HeaderValue::from_static("test-kms-key-id"));
let usecase = DefaultObjectUsecase::without_context();
let fs = FS::new();
let err = usecase.execute_put_object(&fs, req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::NotImplemented);
}
}