refactor: improve code quality with safer error handling, trait decomposition, and dead code cleanup (#1997)

This commit is contained in:
安正超
2026-02-28 01:19:47 +08:00
committed by GitHub
parent 7ce23c6b54
commit af6c32efac
46 changed files with 939 additions and 1612 deletions
+129 -278
View File
@@ -13,12 +13,11 @@
// limitations under the License.
//! Object application use-case contracts.
#![allow(dead_code)]
use crate::app::context::{AppContext, default_notify_interface, get_global_app_context};
use crate::config::workload_profiles::RustFSBufferConfig;
use crate::error::ApiError;
use crate::storage::access::{ReqInfo, authorize_request, has_bypass_governance_header};
use crate::storage::access::{ReqInfo, authorize_request, has_bypass_governance_header, req_info_mut};
use crate::storage::concurrency::{
CachedGetObject, ConcurrencyManager, GetObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager,
};
@@ -38,7 +37,6 @@ use datafusion::arrow::{
use futures::StreamExt;
use http::{HeaderMap, HeaderValue, StatusCode};
use metrics::{counter, histogram};
use rustfs_ecstore::StorageAPI;
use rustfs_ecstore::bucket::quota::checker::QuotaChecker;
use rustfs_ecstore::bucket::{
lifecycle::{
@@ -65,7 +63,8 @@ use rustfs_ecstore::error::{StorageError, is_err_bucket_not_found, is_err_object
use rustfs_ecstore::new_object_layer_fn;
use rustfs_ecstore::set_disk::is_valid_storage_class;
use rustfs_ecstore::store_api::{
BucketOptions, HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOptions, ObjectToDelete, PutObjReader,
BucketOperations, BucketOptions, HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOperations, ObjectOptions, ObjectToDelete,
PutObjReader,
};
use rustfs_filemeta::{
REPLICATE_INCOMING_DELETE, ReplicationStatusType, ReplicationType, RestoreStatusOps, VersionPurgeStatusType,
@@ -112,7 +111,49 @@ use tokio_util::io::{ReaderStream, StreamReader};
use tracing::{debug, error, info, instrument, warn};
use uuid::Uuid;
pub type ObjectUsecaseResult<T> = Result<T, ApiError>;
/// Extract trailing-header checksum values, overriding the corresponding input fields.
fn apply_trailing_checksums(
algorithm: Option<&str>,
trailing_headers: &Option<s3s::TrailingHeaders>,
checksums: &mut PutObjectChecksums,
) {
let Some(alg) = algorithm else { return };
let Some(checksum_str) = trailing_headers.as_ref().and_then(|trailer| {
let key = match alg {
ChecksumAlgorithm::CRC32 => rustfs_rio::ChecksumType::CRC32.key(),
ChecksumAlgorithm::CRC32C => rustfs_rio::ChecksumType::CRC32C.key(),
ChecksumAlgorithm::SHA1 => rustfs_rio::ChecksumType::SHA1.key(),
ChecksumAlgorithm::SHA256 => rustfs_rio::ChecksumType::SHA256.key(),
ChecksumAlgorithm::CRC64NVME => rustfs_rio::ChecksumType::CRC64_NVME.key(),
_ => return None,
};
trailer.read(|headers| {
headers
.get(key.unwrap_or_default())
.and_then(|value| value.to_str().ok().map(|s| s.to_string()))
})
}) else {
return;
};
match alg {
ChecksumAlgorithm::CRC32 => checksums.crc32 = checksum_str,
ChecksumAlgorithm::CRC32C => checksums.crc32c = checksum_str,
ChecksumAlgorithm::SHA1 => checksums.sha1 = checksum_str,
ChecksumAlgorithm::SHA256 => checksums.sha256 = checksum_str,
ChecksumAlgorithm::CRC64NVME => checksums.crc64nvme = checksum_str,
_ => (),
}
}
#[derive(Default)]
struct PutObjectChecksums {
crc32: Option<String>,
crc32c: Option<String>,
sha1: Option<String>,
sha256: Option<String>,
crc64nvme: Option<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());
@@ -129,67 +170,13 @@ fn normalize_delete_objects_version_id(version_id: Option<String>) -> Result<(Op
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PutObjectRequest {
pub bucket: String,
pub key: String,
pub content_length: Option<i64>,
pub content_type: Option<String>,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PutObjectResponse {
pub etag: String,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GetObjectRequest {
pub bucket: String,
pub key: String,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct GetObjectResponse {
pub etag: Option<String>,
pub content_length: i64,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeleteObjectRequest {
pub bucket: String,
pub key: String,
pub version_id: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct DeleteObjectResponse {
pub delete_marker: Option<bool>,
pub version_id: Option<String>,
}
#[async_trait::async_trait]
pub trait ObjectUsecase: Send + Sync {
async fn put_object(&self, req: PutObjectRequest) -> ObjectUsecaseResult<PutObjectResponse>;
async fn get_object(&self, req: GetObjectRequest) -> ObjectUsecaseResult<GetObjectResponse>;
async fn delete_object(&self, req: DeleteObjectRequest) -> ObjectUsecaseResult<DeleteObjectResponse>;
}
#[derive(Clone, Default)]
pub struct DefaultObjectUsecase {
context: Option<Arc<AppContext>>,
}
impl DefaultObjectUsecase {
pub fn new(context: Arc<AppContext>) -> Self {
Self { context: Some(context) }
}
#[cfg(test)]
pub fn without_context() -> Self {
Self { context: None }
}
@@ -200,10 +187,6 @@ impl DefaultObjectUsecase {
}
}
pub fn context(&self) -> Option<Arc<AppContext>> {
self.context.clone()
}
fn bucket_metadata_sys(&self) -> Option<Arc<RwLock<metadata_sys::BucketMetadataSys>>> {
self.context.as_ref().and_then(|context| context.bucket_metadata().handle())
}
@@ -216,6 +199,35 @@ impl DefaultObjectUsecase {
.unwrap_or_else(|| RustFSBufferConfig::default().base_config.default_unknown)
}
async fn check_bucket_quota(&self, bucket: &str, op: QuotaOperation, size: u64) -> S3Result<()> {
let Some(metadata_sys) = self.bucket_metadata_sys() else {
return Ok(());
};
let quota_checker = QuotaChecker::new(metadata_sys);
match quota_checker.check_quota(bucket, op, size).await {
Ok(result) if !result.allowed => Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!(
"Bucket quota exceeded. Current usage: {} bytes, limit: {} bytes",
result.current_usage.unwrap_or(0),
result.quota_limit.unwrap_or(0)
),
)),
Err(e) => {
warn!("Quota check failed for bucket {bucket}: {e}, allowing operation");
Ok(())
}
_ => Ok(()),
}
}
fn spawn_cache_invalidation(bucket: String, key: String, version_id: Option<String>) {
let manager = get_concurrency_manager();
tokio::spawn(async move {
manager.invalidate_cache_versioned(&bucket, &key, version_id.as_deref()).await;
});
}
#[instrument(level = "debug", skip(self, _fs, req))]
pub async fn execute_put_object(&self, _fs: &FS, req: S3Request<PutObjectInput>) -> S3Result<S3Response<PutObjectOutput>> {
if let Some(context) = &self.context {
@@ -273,32 +285,9 @@ impl DefaultObjectUsecase {
return Err(s3_error!(AccessDenied, "Access Denied"));
}
// check quota for put operation
if let Some(size) = content_length
&& let Some(metadata_sys) = self.bucket_metadata_sys()
{
let quota_checker = QuotaChecker::new(metadata_sys);
match quota_checker
.check_quota(&bucket, QuotaOperation::PutObject, size as u64)
.await
{
Ok(check_result) => {
if !check_result.allowed {
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!(
"Bucket quota exceeded. Current usage: {} bytes, limit: {} bytes",
check_result.current_usage.unwrap_or(0),
check_result.quota_limit.unwrap_or(0)
),
));
}
}
Err(e) => {
warn!("Quota check failed for bucket {}: {}, allowing operation", bucket, e);
}
}
if let Some(size) = content_length {
self.check_bucket_quota(&bucket, QuotaOperation::PutObject, size as u64)
.await?;
}
let Some(body) = body else { return Err(s3_error!(IncompleteBody)) };
@@ -330,10 +319,6 @@ impl DefaultObjectUsecase {
StreamReader::new(body.map(|f| f.map_err(|e| std::io::Error::other(e.to_string())))),
);
// let body = Box::new(StreamReader::new(body.map(|f| f.map_err(|e| std::io::Error::other(e.to_string())))));
// let mut reader = PutObjReader::new(body, content_length as usize);
let store = get_validated_store(&bucket).await?;
// TDD: Get bucket default encryption configuration
@@ -514,10 +499,6 @@ impl DefaultObjectUsecase {
// Fast in-memory update for immediate quota consistency
rustfs_ecstore::data_usage::increment_bucket_usage_memory(&bucket, obj_info.size as u64).await;
// Invalidate cache for the written object to prevent stale data
let manager = get_concurrency_manager();
let put_bucket = bucket.clone();
let put_key = key.clone();
let put_version = obj_info.version_id.map(|v| v.to_string());
helper = helper.object(obj_info.clone());
@@ -525,12 +506,7 @@ impl DefaultObjectUsecase {
helper = helper.version_id(version_id.clone());
}
let put_version_clone = put_version.clone();
tokio::spawn(async move {
manager
.invalidate_cache_versioned(&put_bucket, &put_key, put_version_clone.as_deref())
.await;
});
Self::spawn_cache_invalidation(bucket.clone(), key.clone(), put_version.clone());
let e_tag = obj_info.etag.clone().map(|etag| to_s3s_etag(&etag));
@@ -543,77 +519,35 @@ impl DefaultObjectUsecase {
schedule_replication(obj_info, store, dsc, ReplicationType::Object).await;
}
let mut checksum_crc32 = input.checksum_crc32;
let mut checksum_crc32c = input.checksum_crc32c;
let mut checksum_sha1 = input.checksum_sha1;
let mut checksum_sha256 = input.checksum_sha256;
let mut checksum_crc64nvme = input.checksum_crc64nvme;
if let Some(alg) = &input.checksum_algorithm
&& let Some(Some(checksum_str)) = req.trailing_headers.as_ref().map(|trailer| {
let key = match alg.as_str() {
ChecksumAlgorithm::CRC32 => rustfs_rio::ChecksumType::CRC32.key(),
ChecksumAlgorithm::CRC32C => rustfs_rio::ChecksumType::CRC32C.key(),
ChecksumAlgorithm::SHA1 => rustfs_rio::ChecksumType::SHA1.key(),
ChecksumAlgorithm::SHA256 => rustfs_rio::ChecksumType::SHA256.key(),
ChecksumAlgorithm::CRC64NVME => rustfs_rio::ChecksumType::CRC64_NVME.key(),
_ => return None,
};
trailer.read(|headers| {
headers
.get(key.unwrap_or_default())
.and_then(|value| value.to_str().ok().map(|s| s.to_string()))
})
})
{
match alg.as_str() {
ChecksumAlgorithm::CRC32 => checksum_crc32 = checksum_str,
ChecksumAlgorithm::CRC32C => checksum_crc32c = checksum_str,
ChecksumAlgorithm::SHA1 => checksum_sha1 = checksum_str,
ChecksumAlgorithm::SHA256 => checksum_sha256 = checksum_str,
ChecksumAlgorithm::CRC64NVME => checksum_crc64nvme = checksum_str,
_ => (),
}
}
let mut checksums = PutObjectChecksums {
crc32: input.checksum_crc32,
crc32c: input.checksum_crc32c,
sha1: input.checksum_sha1,
sha256: input.checksum_sha256,
crc64nvme: input.checksum_crc64nvme,
};
apply_trailing_checksums(
input.checksum_algorithm.as_ref().map(|a| a.as_str()),
&req.trailing_headers,
&mut checksums,
);
let output = PutObjectOutput {
e_tag,
server_side_encryption: effective_sse, // TDD: Return effective encryption config
server_side_encryption: effective_sse,
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key_md5: sse_customer_key_md5.clone(),
ssekms_key_id: effective_kms_key_id, // TDD: Return effective KMS key ID
checksum_crc32,
checksum_crc32c,
checksum_sha1,
checksum_sha256,
checksum_crc64nvme,
ssekms_key_id: effective_kms_key_id,
checksum_crc32: checksums.crc32,
checksum_crc32c: checksums.crc32c,
checksum_sha1: checksums.sha1,
checksum_sha256: checksums.sha256,
checksum_crc64nvme: checksums.crc64nvme,
version_id: put_version,
..Default::default()
};
// TODO fix response for POST Policy (multipart/form-data) wait s3s crate update,fix issue #1564
// // If it is a POST Policy(multipart/form-data) path, the PutObjectInput carries the success_action_* field
// // Here, the response is uniformly rewritten, with the default being 204, redirect prioritizing 303, and status supporting 200/201/204
// if input.success_action_status.is_some() || input.success_action_redirect.is_some() {
// let mut form_fields = HashMap::<String, String>::new();
// if let Some(v) = &input.success_action_status {
// form_fields.insert("success_action_status".to_string(), v.to_string());
// }
// if let Some(v) = &input.success_action_redirect {
// form_fields.insert("success_action_redirect".to_string(), v.to_string());
// }
//
// // obj_info.etag has been converted to e_tag (s3s etag) above, so try to pass the original string here
// let etag_str = e_tag.as_ref().map(|v| v.as_str());
//
// // Returns using POST semantics: 204/303/201/200
// let resp = build_post_object_success_response(&form_fields, &bucket, &key, etag_str, None)?;
//
// // Keep helper event complete (note: (StatusCode, Body) is returned here instead of PutObjectOutput)
// let result = Ok(resp);
// let _ = helper.complete(&result);
// return result;
// }
// TODO fix response for POST Policy (multipart/form-data), wait s3s crate update, fix issue #1564
let result = Ok(S3Response::new(output));
let _ = helper.complete(&result);
@@ -1200,7 +1134,10 @@ impl DefaultObjectUsecase {
// based on the wait time. Longer wait times indicate higher system
// load, which triggers more conservative I/O parameters.
let permit_wait_start = std::time::Instant::now();
let _disk_permit = manager.acquire_disk_read_permit().await;
let _disk_permit = manager
.acquire_disk_read_permit()
.await
.map_err(|_| s3_error!(InternalError, "disk read semaphore closed"))?;
let permit_wait_duration = permit_wait_start.elapsed();
// Calculate adaptive I/O strategy from permit wait time
@@ -2107,34 +2044,9 @@ impl DefaultObjectUsecase {
src_info.user_defined.insert(k, v);
}
// check quota for copy operation
let has_bucket_metadata = if let Some(metadata_sys) = self.bucket_metadata_sys() {
let quota_checker = QuotaChecker::new(metadata_sys);
match quota_checker
.check_quota(&bucket, QuotaOperation::CopyObject, src_info.size as u64)
.await
{
Ok(check_result) => {
if !check_result.allowed {
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!(
"Bucket quota exceeded. Current usage: {} bytes, limit: {} bytes",
check_result.current_usage.unwrap_or(0),
check_result.quota_limit.unwrap_or(0)
),
));
}
}
Err(e) => {
warn!("Quota check failed for bucket {}: {}, allowing operation", bucket, e);
}
}
true
} else {
false
};
self.check_bucket_quota(&bucket, QuotaOperation::CopyObject, src_info.size as u64)
.await?;
let has_bucket_metadata = self.bucket_metadata_sys().is_some();
let oi = store
.copy_object(&src_bucket, &src_key, &bucket, &key, &mut src_info, &src_opts, &dst_opts)
@@ -2146,17 +2058,8 @@ impl DefaultObjectUsecase {
rustfs_ecstore::data_usage::increment_bucket_usage_memory(&bucket, oi.size as u64).await;
}
// Invalidate cache for the destination object to prevent stale data
let manager = get_concurrency_manager();
let dest_bucket = bucket.clone();
let dest_key = key.clone();
let dest_version = oi.version_id.map(|v| v.to_string());
let dest_version_clone = dest_version.clone();
tokio::spawn(async move {
manager
.invalidate_cache_versioned(&dest_bucket, &dest_key, dest_version_clone.as_deref())
.await;
});
Self::spawn_cache_invalidation(bucket.clone(), key.clone(), dest_version.clone());
// warn!("copy_object oi {:?}", &oi);
let object_info = oi.clone();
@@ -2254,7 +2157,7 @@ impl DefaultObjectUsecase {
};
{
let req_info = req.extensions.get_mut::<ReqInfo>().expect("ReqInfo not found");
let req_info = req_info_mut(&mut req)?;
req_info.bucket = Some(bucket.clone());
req_info.object = Some(obj_id.key.clone());
req_info.version_id = version_id.clone();
@@ -2606,16 +2509,7 @@ impl DefaultObjectUsecase {
// Fast in-memory update for immediate quota consistency
rustfs_ecstore::data_usage::decrement_bucket_usage_memory(&bucket, obj_info.size as u64).await;
// Invalidate cache for the deleted object
let manager = get_concurrency_manager();
let del_bucket = bucket.clone();
let del_key = key.clone();
let del_version = obj_info.version_id.map(|v| v.to_string());
tokio::spawn(async move {
manager
.invalidate_cache_versioned(&del_bucket, &del_key, del_version.as_deref())
.await;
});
Self::spawn_cache_invalidation(bucket.clone(), key.clone(), obj_info.version_id.map(|v| v.to_string()));
if obj_info.name.is_empty() {
if replicate_force_delete {
@@ -3489,49 +3383,30 @@ impl DefaultObjectUsecase {
}
}
let mut checksum_crc32 = input.checksum_crc32;
let mut checksum_crc32c = input.checksum_crc32c;
let mut checksum_sha1 = input.checksum_sha1;
let mut checksum_sha256 = input.checksum_sha256;
let mut checksum_crc64nvme = input.checksum_crc64nvme;
if let Some(alg) = &input.checksum_algorithm
&& let Some(Some(checksum_str)) = req.trailing_headers.as_ref().map(|trailer| {
let key = match alg.as_str() {
ChecksumAlgorithm::CRC32 => rustfs_rio::ChecksumType::CRC32.key(),
ChecksumAlgorithm::CRC32C => rustfs_rio::ChecksumType::CRC32C.key(),
ChecksumAlgorithm::SHA1 => rustfs_rio::ChecksumType::SHA1.key(),
ChecksumAlgorithm::SHA256 => rustfs_rio::ChecksumType::SHA256.key(),
ChecksumAlgorithm::CRC64NVME => rustfs_rio::ChecksumType::CRC64_NVME.key(),
_ => return None,
};
trailer.read(|headers| {
headers
.get(key.unwrap_or_default())
.and_then(|value| value.to_str().ok().map(|s| s.to_string()))
})
})
{
match alg.as_str() {
ChecksumAlgorithm::CRC32 => checksum_crc32 = checksum_str,
ChecksumAlgorithm::CRC32C => checksum_crc32c = checksum_str,
ChecksumAlgorithm::SHA1 => checksum_sha1 = checksum_str,
ChecksumAlgorithm::SHA256 => checksum_sha256 = checksum_str,
ChecksumAlgorithm::CRC64NVME => checksum_crc64nvme = checksum_str,
_ => (),
}
}
let mut checksums = PutObjectChecksums {
crc32: input.checksum_crc32,
crc32c: input.checksum_crc32c,
sha1: input.checksum_sha1,
sha256: input.checksum_sha256,
crc64nvme: input.checksum_crc64nvme,
};
apply_trailing_checksums(
input.checksum_algorithm.as_ref().map(|a| a.as_str()),
&req.trailing_headers,
&mut checksums,
);
warn!(
"put object extract checksum_crc32={checksum_crc32:?}, checksum_crc32c={checksum_crc32c:?}, checksum_sha1={checksum_sha1:?}, checksum_sha256={checksum_sha256:?}, checksum_crc64nvme={checksum_crc64nvme:?}",
"put object extract checksum_crc32={:?}, checksum_crc32c={:?}, checksum_sha1={:?}, checksum_sha256={:?}, checksum_crc64nvme={:?}",
checksums.crc32, checksums.crc32c, checksums.sha1, checksums.sha256, checksums.crc64nvme,
);
let output = PutObjectOutput {
checksum_crc32,
checksum_crc32c,
checksum_sha1,
checksum_sha256,
checksum_crc64nvme,
checksum_crc32: checksums.crc32,
checksum_crc32c: checksums.crc32c,
checksum_sha1: checksums.sha1,
checksum_sha256: checksums.sha256,
checksum_crc64nvme: checksums.crc64nvme,
..Default::default()
};
let result = Ok(S3Response::new(output));
@@ -3540,30 +3415,6 @@ impl DefaultObjectUsecase {
}
}
#[async_trait::async_trait]
impl ObjectUsecase for DefaultObjectUsecase {
async fn put_object(&self, req: PutObjectRequest) -> ObjectUsecaseResult<PutObjectResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultObjectUsecase::put_object DTO path is not implemented yet",
)))
}
async fn get_object(&self, req: GetObjectRequest) -> ObjectUsecaseResult<GetObjectResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultObjectUsecase::get_object DTO path is not implemented yet",
)))
}
async fn delete_object(&self, req: DeleteObjectRequest) -> ObjectUsecaseResult<DeleteObjectResponse> {
let _ = req;
Err(ApiError::from(StorageError::other(
"DefaultObjectUsecase::delete_object DTO path is not implemented yet",
)))
}
}
#[cfg(test)]
mod tests {
use super::*;