refactor(app): route put/get/listv2 through usecases (#1910)

This commit is contained in:
安正超
2026-02-22 23:36:48 +08:00
committed by GitHub
parent 84053484e6
commit cf1d109bb9
6 changed files with 1964 additions and 1306 deletions
+67
View File
@@ -0,0 +1,67 @@
// 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.
//! Admin application use-case contracts.
#![allow(dead_code)]
use crate::app::context::AppContext;
use crate::error::ApiError;
use rustfs_ecstore::admin_server_info::get_server_info;
use rustfs_madmin::InfoMessage;
use std::sync::Arc;
pub type AdminUsecaseResult<T> = Result<T, ApiError>;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct QueryServerInfoRequest {
pub include_pools: bool,
}
pub struct QueryServerInfoResponse {
pub info: InfoMessage,
}
impl std::fmt::Debug for QueryServerInfoResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("QueryServerInfoResponse").finish_non_exhaustive()
}
}
#[async_trait::async_trait]
pub trait AdminUsecase: Send + Sync {
async fn query_server_info(&self, req: QueryServerInfoRequest) -> AdminUsecaseResult<QueryServerInfoResponse>;
}
#[derive(Clone)]
pub struct DefaultAdminUsecase {
context: Arc<AppContext>,
}
impl DefaultAdminUsecase {
pub fn new(context: Arc<AppContext>) -> Self {
Self { context }
}
pub fn context(&self) -> Arc<AppContext> {
self.context.clone()
}
}
#[async_trait::async_trait]
impl AdminUsecase for DefaultAdminUsecase {
async fn query_server_info(&self, req: QueryServerInfoRequest) -> AdminUsecaseResult<QueryServerInfoResponse> {
let info = get_server_info(req.include_pools).await;
Ok(QueryServerInfoResponse { info })
}
}
+343
View File
@@ -0,0 +1,343 @@
// 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.
//! Bucket application use-case contracts.
#![allow(dead_code)]
use crate::app::context::{AppContext, get_global_app_context};
use crate::error::ApiError;
use crate::storage::*;
use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
use rustfs_ecstore::error::StorageError;
use rustfs_ecstore::store_api::StorageAPI;
use s3s::dto::*;
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
use std::sync::Arc;
use tracing::{debug, instrument};
use urlencoding::encode;
pub type BucketUsecaseResult<T> = Result<T, ApiError>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CreateBucketRequest {
pub bucket: String,
pub object_lock_enabled: Option<bool>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct CreateBucketResponse;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeleteBucketRequest {
pub bucket: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct DeleteBucketResponse;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HeadBucketRequest {
pub bucket: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct HeadBucketResponse;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ListObjectsV2Request {
pub bucket: String,
pub prefix: Option<String>,
pub delimiter: Option<String>,
pub continuation_token: Option<String>,
pub max_keys: Option<i32>,
pub fetch_owner: Option<bool>,
pub start_after: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ListObjectsV2Item {
pub key: String,
pub etag: Option<String>,
pub size: i64,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ListObjectsV2Response {
pub objects: Vec<ListObjectsV2Item>,
pub common_prefixes: Vec<String>,
pub key_count: i32,
pub is_truncated: bool,
pub next_continuation_token: Option<String>,
}
#[async_trait::async_trait]
pub trait BucketUsecase: Send + Sync {
async fn create_bucket(&self, req: CreateBucketRequest) -> BucketUsecaseResult<CreateBucketResponse>;
async fn delete_bucket(&self, req: DeleteBucketRequest) -> BucketUsecaseResult<DeleteBucketResponse>;
async fn head_bucket(&self, req: HeadBucketRequest) -> BucketUsecaseResult<HeadBucketResponse>;
async fn list_objects_v2(&self, req: ListObjectsV2Request) -> BucketUsecaseResult<ListObjectsV2Response>;
}
#[derive(Clone, Default)]
pub struct DefaultBucketUsecase {
context: Option<Arc<AppContext>>,
}
impl DefaultBucketUsecase {
pub fn new(context: Arc<AppContext>) -> Self {
Self { context: Some(context) }
}
pub fn without_context() -> Self {
Self { context: None }
}
pub fn from_global() -> Self {
Self {
context: get_global_app_context(),
}
}
pub fn context(&self) -> Option<Arc<AppContext>> {
self.context.clone()
}
#[instrument(level = "debug", skip(self, req))]
pub async fn execute_list_objects_v2(&self, req: S3Request<ListObjectsV2Input>) -> S3Result<S3Response<ListObjectsV2Output>> {
if let Some(context) = &self.context {
let _ = context.object_store();
}
// warn!("list_objects_v2 req {:?}", &req.input);
let ListObjectsV2Input {
bucket,
continuation_token,
delimiter,
encoding_type,
fetch_owner,
max_keys,
prefix,
start_after,
..
} = req.input;
let prefix = prefix.unwrap_or_default();
// Log debug info for prefixes with special characters to help diagnose encoding issues
if prefix.contains([' ', '+', '%', '\n', '\r', '\0']) {
debug!("LIST objects with special characters in prefix: {:?}", prefix);
}
let max_keys = max_keys.unwrap_or(1000);
if max_keys < 0 {
return Err(S3Error::with_message(S3ErrorCode::InvalidArgument, "Invalid max keys".to_string()));
}
let delimiter = delimiter.filter(|v| !v.is_empty());
validate_list_object_unordered_with_delimiter(delimiter.as_ref(), req.uri.query())?;
// Save original start_after for response (per S3 API spec, must echo back if provided)
let response_start_after = start_after.clone();
let start_after_for_query = start_after.filter(|v| !v.is_empty());
// Save original continuation_token for response (per S3 API spec, must echo back if provided)
// Note: empty string should still be echoed back in the response
let response_continuation_token = continuation_token.clone();
let continuation_token_for_query = continuation_token.filter(|v| !v.is_empty());
// Decode continuation_token from base64 for internal use
let decoded_continuation_token = continuation_token_for_query
.map(|token| {
base64_simd::STANDARD
.decode_to_vec(token.as_bytes())
.map_err(|_| s3_error!(InvalidArgument, "Invalid continuation token"))
.and_then(|bytes| {
String::from_utf8(bytes).map_err(|_| s3_error!(InvalidArgument, "Invalid continuation token"))
})
})
.transpose()?;
let store = get_validated_store(&bucket).await?;
let incl_deleted = req
.headers
.get(rustfs_utils::http::headers::RUSTFS_INCLUDE_DELETED)
.is_some_and(|v| v.to_str().unwrap_or_default() == "true");
let object_infos = store
.list_objects_v2(
&bucket,
&prefix,
decoded_continuation_token,
delimiter.clone(),
max_keys,
fetch_owner.unwrap_or_default(),
start_after_for_query,
incl_deleted,
)
.await
.map_err(ApiError::from)?;
// warn!("object_infos objects {:?}", object_infos.objects);
// Apply URL encoding if encoding_type is "url"
// Note: S3 URL encoding should encode special characters but preserve path separators (/)
let should_encode = encoding_type.as_ref().map(|e| e.as_str() == "url").unwrap_or(false);
// Helper function to encode S3 keys/prefixes (preserving /)
// S3 URL encoding encodes special characters but keeps '/' unencoded
let encode_s3_name = |name: &str| -> String {
name.split('/')
.map(|part| encode(part).to_string())
.collect::<Vec<_>>()
.join("/")
};
let objects: Vec<Object> = object_infos
.objects
.iter()
.filter(|v| !v.name.is_empty())
.map(|v| {
let key = if should_encode {
encode_s3_name(&v.name)
} else {
v.name.to_owned()
};
let mut obj = Object {
key: Some(key),
last_modified: v.mod_time.map(Timestamp::from),
size: Some(v.get_actual_size().unwrap_or_default()),
e_tag: v.etag.clone().map(|etag| to_s3s_etag(&etag)),
storage_class: v.storage_class.clone().map(ObjectStorageClass::from),
..Default::default()
};
if fetch_owner.is_some_and(|v| v) {
obj.owner = Some(Owner {
display_name: Some("rustfs".to_owned()),
id: Some("v0.1".to_owned()),
});
}
obj
})
.collect();
let common_prefixes: Vec<CommonPrefix> = object_infos
.prefixes
.into_iter()
.map(|v| {
let prefix = if should_encode { encode_s3_name(&v) } else { v };
CommonPrefix { prefix: Some(prefix) }
})
.collect();
// KeyCount should include both objects and common prefixes per S3 API spec
let key_count = (objects.len() + common_prefixes.len()) as i32;
// Encode next_continuation_token to base64
let next_continuation_token = object_infos
.next_continuation_token
.map(|token| base64_simd::STANDARD.encode_to_string(token.as_bytes()));
let output = ListObjectsV2Output {
is_truncated: Some(object_infos.is_truncated),
continuation_token: response_continuation_token,
next_continuation_token,
start_after: response_start_after,
key_count: Some(key_count),
max_keys: Some(max_keys),
contents: Some(objects),
delimiter,
encoding_type: encoding_type.clone(),
name: Some(bucket),
prefix: Some(prefix),
common_prefixes: Some(common_prefixes),
..Default::default()
};
// let output = ListObjectsV2Output { ..Default::default() };
Ok(S3Response::new(output))
}
}
#[async_trait::async_trait]
impl BucketUsecase for DefaultBucketUsecase {
async fn create_bucket(&self, req: CreateBucketRequest) -> BucketUsecaseResult<CreateBucketResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultBucketUsecase::create_bucket DTO path is not implemented yet",
)))
}
async fn delete_bucket(&self, req: DeleteBucketRequest) -> BucketUsecaseResult<DeleteBucketResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultBucketUsecase::delete_bucket DTO path is not implemented yet",
)))
}
async fn head_bucket(&self, req: HeadBucketRequest) -> BucketUsecaseResult<HeadBucketResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultBucketUsecase::head_bucket DTO path is not implemented yet",
)))
}
async fn list_objects_v2(&self, req: ListObjectsV2Request) -> BucketUsecaseResult<ListObjectsV2Response> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultBucketUsecase::list_objects_v2 DTO path is not implemented yet",
)))
}
}
#[cfg(test)]
mod tests {
use super::*;
use http::{Extensions, HeaderMap, Method, Uri};
fn build_request<T>(input: T, method: Method) -> S3Request<T> {
S3Request {
input,
method,
uri: Uri::from_static("/"),
headers: HeaderMap::new(),
extensions: Extensions::new(),
credentials: None,
region: None,
service: None,
trailing_headers: None,
}
}
#[tokio::test]
async fn execute_list_objects_v2_rejects_negative_max_keys() {
let input = ListObjectsV2Input::builder()
.bucket("test-bucket".to_string())
.max_keys(Some(-1))
.build()
.unwrap();
let req = build_request(input, Method::GET);
let usecase = DefaultBucketUsecase::without_context();
let err = usecase.execute_list_objects_v2(req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
}
}
+4
View File
@@ -15,4 +15,8 @@
//! Application layer module entry.
//! Concrete use-case modules will be introduced incrementally in Phase 3.
pub mod admin_usecase;
pub mod bucket_usecase;
pub mod context;
pub mod multipart_usecase;
pub mod object_usecase;
+156
View File
@@ -0,0 +1,156 @@
// 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.
//! Multipart application use-case contracts.
#![allow(dead_code)]
use crate::app::context::AppContext;
use crate::error::ApiError;
use rustfs_ecstore::error::StorageError;
use std::collections::HashMap;
use std::sync::Arc;
pub type MultipartUsecaseResult<T> = Result<T, ApiError>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CreateMultipartUploadRequest {
pub bucket: String,
pub key: String,
pub metadata: HashMap<String, String>,
pub content_type: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CreateMultipartUploadResponse {
pub upload_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UploadPartRequest {
pub bucket: String,
pub key: String,
pub upload_id: String,
pub part_number: i32,
pub content_length: Option<i64>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UploadPartResponse {
pub etag: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompleteMultipartUploadPart {
pub part_number: i32,
pub etag: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompleteMultipartUploadRequest {
pub bucket: String,
pub key: String,
pub upload_id: String,
pub parts: Vec<CompleteMultipartUploadPart>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct CompleteMultipartUploadResponse {
pub etag: Option<String>,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AbortMultipartUploadRequest {
pub bucket: String,
pub key: String,
pub upload_id: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct AbortMultipartUploadResponse;
#[async_trait::async_trait]
pub trait MultipartUsecase: Send + Sync {
async fn create_multipart_upload(
&self,
req: CreateMultipartUploadRequest,
) -> MultipartUsecaseResult<CreateMultipartUploadResponse>;
async fn upload_part(&self, req: UploadPartRequest) -> MultipartUsecaseResult<UploadPartResponse>;
async fn complete_multipart_upload(
&self,
req: CompleteMultipartUploadRequest,
) -> MultipartUsecaseResult<CompleteMultipartUploadResponse>;
async fn abort_multipart_upload(
&self,
req: AbortMultipartUploadRequest,
) -> MultipartUsecaseResult<AbortMultipartUploadResponse>;
}
#[derive(Clone)]
pub struct DefaultMultipartUsecase {
context: Arc<AppContext>,
}
impl DefaultMultipartUsecase {
pub fn new(context: Arc<AppContext>) -> Self {
Self { context }
}
pub fn context(&self) -> Arc<AppContext> {
self.context.clone()
}
}
#[async_trait::async_trait]
impl MultipartUsecase for DefaultMultipartUsecase {
async fn create_multipart_upload(
&self,
req: CreateMultipartUploadRequest,
) -> MultipartUsecaseResult<CreateMultipartUploadResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultMultipartUsecase::create_multipart_upload is not implemented yet",
)))
}
async fn upload_part(&self, req: UploadPartRequest) -> MultipartUsecaseResult<UploadPartResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultMultipartUsecase::upload_part is not implemented yet",
)))
}
async fn complete_multipart_upload(
&self,
req: CompleteMultipartUploadRequest,
) -> MultipartUsecaseResult<CompleteMultipartUploadResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultMultipartUsecase::complete_multipart_upload is not implemented yet",
)))
}
async fn abort_multipart_upload(
&self,
req: AbortMultipartUploadRequest,
) -> MultipartUsecaseResult<AbortMultipartUploadResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultMultipartUsecase::abort_multipart_upload is not implemented yet",
)))
}
}
File diff suppressed because it is too large Load Diff
+23 -1306
View File
File diff suppressed because it is too large Load Diff