mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 12:26:37 +00:00
refactor(app): reuse multipart uploads helpers (#2447)
This commit is contained in:
@@ -25,7 +25,10 @@ use crate::storage::options::{
|
|||||||
copy_src_opts, extract_metadata, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts,
|
copy_src_opts, extract_metadata, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts,
|
||||||
parse_copy_source_range, put_opts, validate_archive_content_encoding,
|
parse_copy_source_range, put_opts, validate_archive_content_encoding,
|
||||||
};
|
};
|
||||||
use crate::storage::s3_api::multipart::{build_list_parts_output, parse_list_parts_params};
|
use crate::storage::s3_api::multipart::{
|
||||||
|
ListMultipartUploadsParams, build_list_multipart_uploads_output, build_list_parts_output,
|
||||||
|
parse_list_multipart_uploads_params, parse_list_parts_params,
|
||||||
|
};
|
||||||
use crate::storage::*;
|
use crate::storage::*;
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use futures::StreamExt;
|
use futures::StreamExt;
|
||||||
@@ -42,7 +45,7 @@ use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
|
|||||||
use rustfs_ecstore::compress::is_compressible;
|
use rustfs_ecstore::compress::is_compressible;
|
||||||
use rustfs_ecstore::error::{StorageError, is_err_object_not_found, is_err_version_not_found};
|
use rustfs_ecstore::error::{StorageError, is_err_object_not_found, is_err_version_not_found};
|
||||||
use rustfs_ecstore::new_object_layer_fn;
|
use rustfs_ecstore::new_object_layer_fn;
|
||||||
use rustfs_ecstore::set_disk::{MAX_PARTS_COUNT, is_valid_storage_class};
|
use rustfs_ecstore::set_disk::is_valid_storage_class;
|
||||||
use rustfs_ecstore::store_api::{
|
use rustfs_ecstore::store_api::{
|
||||||
ChunkNativePutData, CompletePart, HTTPRangeSpec, MultipartUploadResult, ObjectIO, ObjectOptions,
|
ChunkNativePutData, CompletePart, HTTPRangeSpec, MultipartUploadResult, ObjectIO, ObjectOptions,
|
||||||
};
|
};
|
||||||
@@ -917,55 +920,22 @@ impl DefaultMultipartUsecase {
|
|||||||
..
|
..
|
||||||
} = req.input;
|
} = req.input;
|
||||||
|
|
||||||
|
let ListMultipartUploadsParams {
|
||||||
|
prefix,
|
||||||
|
key_marker,
|
||||||
|
max_uploads,
|
||||||
|
} = parse_list_multipart_uploads_params(prefix, key_marker, max_uploads)?;
|
||||||
|
|
||||||
let Some(store) = new_object_layer_fn() else {
|
let Some(store) = new_object_layer_fn() else {
|
||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||||
};
|
};
|
||||||
|
|
||||||
let prefix = prefix.unwrap_or_default();
|
|
||||||
let max_uploads = max_uploads.map(|x| x as usize).unwrap_or(MAX_PARTS_COUNT);
|
|
||||||
|
|
||||||
if let Some(key_marker) = &key_marker
|
|
||||||
&& !key_marker.starts_with(prefix.as_str())
|
|
||||||
{
|
|
||||||
return Err(s3_error!(NotImplemented, "Invalid key marker"));
|
|
||||||
}
|
|
||||||
|
|
||||||
let result = store
|
let result = store
|
||||||
.list_multipart_uploads(&bucket, &prefix, delimiter, key_marker, upload_id_marker, max_uploads)
|
.list_multipart_uploads(&bucket, &prefix, delimiter, key_marker, upload_id_marker, max_uploads)
|
||||||
.await
|
.await
|
||||||
.map_err(ApiError::from)?;
|
.map_err(ApiError::from)?;
|
||||||
|
|
||||||
let output = ListMultipartUploadsOutput {
|
Ok(S3Response::new(build_list_multipart_uploads_output(bucket, prefix, result)))
|
||||||
bucket: Some(bucket),
|
|
||||||
prefix: Some(prefix),
|
|
||||||
delimiter: result.delimiter,
|
|
||||||
key_marker: result.key_marker,
|
|
||||||
upload_id_marker: result.upload_id_marker,
|
|
||||||
max_uploads: Some(result.max_uploads as i32),
|
|
||||||
is_truncated: Some(result.is_truncated),
|
|
||||||
uploads: Some(
|
|
||||||
result
|
|
||||||
.uploads
|
|
||||||
.into_iter()
|
|
||||||
.map(|u| MultipartUpload {
|
|
||||||
key: Some(u.object),
|
|
||||||
upload_id: Some(u.upload_id),
|
|
||||||
initiated: u.initiated.map(Timestamp::from),
|
|
||||||
..Default::default()
|
|
||||||
})
|
|
||||||
.collect(),
|
|
||||||
),
|
|
||||||
common_prefixes: Some(
|
|
||||||
result
|
|
||||||
.common_prefixes
|
|
||||||
.into_iter()
|
|
||||||
.map(|c| CommonPrefix { prefix: Some(c) })
|
|
||||||
.collect(),
|
|
||||||
),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
|
|
||||||
Ok(S3Response::new(output))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn execute_list_parts(&self, req: S3Request<ListPartsInput>) -> S3Result<S3Response<ListPartsOutput>> {
|
pub async fn execute_list_parts(&self, req: S3Request<ListPartsInput>) -> S3Result<S3Response<ListPartsOutput>> {
|
||||||
@@ -1456,6 +1426,36 @@ mod tests {
|
|||||||
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
assert_eq!(err.code(), &S3ErrorCode::InternalError);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn execute_list_multipart_uploads_rejects_invalid_key_marker_before_store_lookup() {
|
||||||
|
let input = ListMultipartUploadsInput::builder()
|
||||||
|
.bucket("bucket".to_string())
|
||||||
|
.prefix(Some("prefix/".to_string()))
|
||||||
|
.key_marker(Some("other/key".to_string()))
|
||||||
|
.build()
|
||||||
|
.unwrap();
|
||||||
|
let req = build_request(input, Method::GET);
|
||||||
|
|
||||||
|
let err = make_usecase().execute_list_multipart_uploads(req).await.unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::NotImplemented);
|
||||||
|
assert_eq!(err.message(), Some("Invalid key marker"));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn execute_list_multipart_uploads_rejects_invalid_max_uploads_before_store_lookup() {
|
||||||
|
let input = ListMultipartUploadsInput::builder()
|
||||||
|
.bucket("bucket".to_string())
|
||||||
|
.max_uploads(Some(0))
|
||||||
|
.build()
|
||||||
|
.unwrap();
|
||||||
|
let req = build_request(input, Method::GET);
|
||||||
|
let expected = "max-uploads must be between 1 and 1000";
|
||||||
|
|
||||||
|
let err = make_usecase().execute_list_multipart_uploads(req).await.unwrap_err();
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
|
||||||
|
assert_eq!(err.message(), Some(expected));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn execute_list_parts_returns_internal_error_when_store_uninitialized() {
|
async fn execute_list_parts_returns_internal_error_when_store_uninitialized() {
|
||||||
let input = ListPartsInput::builder()
|
let input = ListPartsInput::builder()
|
||||||
|
|||||||
@@ -14,11 +14,12 @@
|
|||||||
|
|
||||||
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
|
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
|
||||||
use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
|
use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
|
||||||
use rustfs_ecstore::set_disk::MAX_PARTS_COUNT;
|
|
||||||
use rustfs_ecstore::store_api::{ListMultipartsInfo, ListPartsInfo};
|
use rustfs_ecstore::store_api::{ListMultipartsInfo, ListPartsInfo};
|
||||||
use s3s::dto::{CommonPrefix, ListMultipartUploadsOutput, ListPartsOutput, MultipartUpload, Part, Timestamp};
|
use s3s::dto::{CommonPrefix, ListMultipartUploadsOutput, ListPartsOutput, MultipartUpload, Part, Timestamp};
|
||||||
use s3s::{S3Error, S3ErrorCode};
|
use s3s::{S3Error, S3ErrorCode};
|
||||||
|
|
||||||
|
const MAX_MULTIPART_UPLOADS_LIST: i32 = 1000;
|
||||||
|
|
||||||
#[derive(Debug, PartialEq, Eq)]
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
pub(crate) struct ListPartsParams {
|
pub(crate) struct ListPartsParams {
|
||||||
pub part_number_marker: Option<usize>,
|
pub part_number_marker: Option<usize>,
|
||||||
@@ -110,23 +111,16 @@ pub(crate) fn parse_list_multipart_uploads_params(
|
|||||||
let prefix = prefix.unwrap_or_default();
|
let prefix = prefix.unwrap_or_default();
|
||||||
let max_uploads = match max_uploads {
|
let max_uploads = match max_uploads {
|
||||||
Some(value) => {
|
Some(value) => {
|
||||||
let value = usize::try_from(value).map_err(|_| {
|
if !(1..=MAX_MULTIPART_UPLOADS_LIST).contains(&value) {
|
||||||
S3Error::with_message(
|
|
||||||
S3ErrorCode::InvalidArgument,
|
|
||||||
format!("max-uploads must be between 1 and {}", MAX_PARTS_COUNT),
|
|
||||||
)
|
|
||||||
})?;
|
|
||||||
|
|
||||||
if value == 0 || value > MAX_PARTS_COUNT {
|
|
||||||
return Err(S3Error::with_message(
|
return Err(S3Error::with_message(
|
||||||
S3ErrorCode::InvalidArgument,
|
S3ErrorCode::InvalidArgument,
|
||||||
format!("max-uploads must be between 1 and {}", MAX_PARTS_COUNT),
|
format!("max-uploads must be between 1 and {}", MAX_MULTIPART_UPLOADS_LIST),
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
value
|
value as usize
|
||||||
}
|
}
|
||||||
None => MAX_PARTS_COUNT,
|
None => MAX_MULTIPART_UPLOADS_LIST as usize,
|
||||||
};
|
};
|
||||||
|
|
||||||
if let Some(key_marker) = &key_marker
|
if let Some(key_marker) = &key_marker
|
||||||
@@ -181,12 +175,11 @@ pub(crate) fn build_list_multipart_uploads_output(
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
build_list_multipart_uploads_output, build_list_parts_output, parse_list_multipart_uploads_params,
|
MAX_MULTIPART_UPLOADS_LIST, build_list_multipart_uploads_output, build_list_parts_output,
|
||||||
parse_list_parts_params,
|
parse_list_multipart_uploads_params, parse_list_parts_params,
|
||||||
};
|
};
|
||||||
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
|
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
|
||||||
use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
|
use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
|
||||||
use rustfs_ecstore::set_disk::MAX_PARTS_COUNT;
|
|
||||||
use rustfs_ecstore::store_api::{ListMultipartsInfo, ListPartsInfo, MultipartInfo, PartInfo};
|
use rustfs_ecstore::store_api::{ListMultipartsInfo, ListPartsInfo, MultipartInfo, PartInfo};
|
||||||
use s3s::S3ErrorCode;
|
use s3s::S3ErrorCode;
|
||||||
use s3s::dto::Timestamp;
|
use s3s::dto::Timestamp;
|
||||||
@@ -338,7 +331,7 @@ mod tests {
|
|||||||
let parsed = parse_list_multipart_uploads_params(None, None, None).expect("expected default params");
|
let parsed = parse_list_multipart_uploads_params(None, None, None).expect("expected default params");
|
||||||
assert_eq!(parsed.prefix, "");
|
assert_eq!(parsed.prefix, "");
|
||||||
assert_eq!(parsed.key_marker, None);
|
assert_eq!(parsed.key_marker, None);
|
||||||
assert_eq!(parsed.max_uploads, MAX_PARTS_COUNT);
|
assert_eq!(parsed.max_uploads, MAX_MULTIPART_UPLOADS_LIST as usize);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -360,7 +353,7 @@ mod tests {
|
|||||||
.expect_err("expected invalid max_uploads");
|
.expect_err("expected invalid max_uploads");
|
||||||
assert_eq!(*err.code(), S3ErrorCode::InvalidArgument);
|
assert_eq!(*err.code(), S3ErrorCode::InvalidArgument);
|
||||||
|
|
||||||
let err = parse_list_multipart_uploads_params(Some("prefix/".to_string()), None, Some((MAX_PARTS_COUNT + 1) as i32))
|
let err = parse_list_multipart_uploads_params(Some("prefix/".to_string()), None, Some(MAX_MULTIPART_UPLOADS_LIST + 1))
|
||||||
.expect_err("expected invalid max_uploads");
|
.expect_err("expected invalid max_uploads");
|
||||||
assert_eq!(*err.code(), S3ErrorCode::InvalidArgument);
|
assert_eq!(*err.code(), S3ErrorCode::InvalidArgument);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user