refactor: wrap bucket compat trait methods (#3706)

This commit is contained in:
安正超
2026-06-22 04:32:37 +08:00
committed by GitHub
parent b8a0fc4e04
commit 7eb9a4759e
16 changed files with 206 additions and 46 deletions
@@ -29,6 +29,7 @@ use crate::admin::storage_compat::replication::GLOBAL_REPLICATION_STATS;
use crate::admin::storage_compat::replication::{ResyncOpts, get_global_replication_pool};
use crate::admin::storage_compat::target::{ARN, BucketTarget, BucketTargetType, BucketTargets, Credentials};
use crate::admin::storage_compat::utils::{deserialize, serialize};
use crate::admin::storage_compat::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _};
use crate::admin::storage_compat::{delete_admin_config, read_admin_config, save_admin_config};
use crate::admin::storage_compat::{get_global_deployment_id, get_global_endpoints_opt, get_global_region, global_rustfs_port};
use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body};
@@ -49,8 +50,6 @@ use rustfs_config::{
DEFAULT_CONSOLE_ADDRESS, DEFAULT_DELIMITER, DEFAULT_RUSTFS_TLS_PATH, ENV_RUSTFS_CONSOLE_ADDRESS, ENV_RUSTFS_TLS_PATH,
MAX_ADMIN_REQUEST_BODY_SIZE,
};
use rustfs_ecstore::api::bucket::replication::ReplicationConfigurationExt as _;
use rustfs_ecstore::api::bucket::versioning::VersioningApi as _;
use rustfs_iam::error::is_err_no_such_service_account;
use rustfs_iam::store::{MappedPolicy, UserType};
use rustfs_iam::sys::{
+1 -2
View File
@@ -29,6 +29,7 @@ use crate::admin::storage_compat::replication::{
};
use crate::admin::storage_compat::target::{BucketTarget, BucketTargetType, BucketTargets};
use crate::admin::storage_compat::versioning_sys::BucketVersioningSys;
use crate::admin::storage_compat::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _};
use crate::admin::storage_compat::{get_global_bucket_monitor, get_global_deployment_id, get_global_region};
use crate::app::context::resolve_object_store_handle;
use crate::app::object_usecase::DefaultObjectUsecase;
@@ -60,8 +61,6 @@ use rustfs_config::{
ENABLE_KEY, WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CA, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY, WEBHOOK_ENDPOINT,
WEBHOOK_SKIP_TLS_VERIFY,
};
use rustfs_ecstore::api::bucket::replication::ReplicationConfigurationExt as _;
use rustfs_ecstore::api::bucket::versioning::VersioningApi as _;
use rustfs_filemeta::{ReplicationStatusType, ReplicationType};
use rustfs_madmin::utils::parse_duration;
use rustfs_notify::{Event as NotificationEvent, notification_system};
+29
View File
@@ -55,6 +55,35 @@ pub(crate) type RebalStatus = crate::admin::storage_compat::ecstore_rebalance::R
#[cfg(test)]
pub(crate) type RebalanceInfo = crate::admin::storage_compat::ecstore_rebalance::RebalanceInfo;
pub(crate) trait AdminReplicationConfigExt {
fn filter_target_arns(&self, obj: &replication::ObjectOpts) -> Vec<String>;
fn has_existing_object_replication(&self, arn: &str) -> (bool, bool);
}
impl AdminReplicationConfigExt for s3s::dto::ReplicationConfiguration {
fn filter_target_arns(&self, obj: &replication::ObjectOpts) -> Vec<String> {
<s3s::dto::ReplicationConfiguration as ecstore_bucket::replication::ReplicationConfigurationExt>::filter_target_arns(
self, obj,
)
}
fn has_existing_object_replication(&self, arn: &str) -> (bool, bool) {
<s3s::dto::ReplicationConfiguration as ecstore_bucket::replication::ReplicationConfigurationExt>::has_existing_object_replication(
self, arn,
)
}
}
pub(crate) trait AdminVersioningConfigExt {
fn enabled(&self) -> bool;
}
impl AdminVersioningConfigExt for s3s::dto::VersioningConfiguration {
fn enabled(&self) -> bool {
<s3s::dto::VersioningConfiguration as ecstore_bucket::versioning::VersioningApi>::enabled(self)
}
}
pub(crate) mod bandwidth {
pub(crate) mod monitor {
pub(crate) type BandwidthDetails = crate::admin::storage_compat::ecstore_bucket::bandwidth::monitor::BandwidthDetails;
+1 -2
View File
@@ -24,6 +24,7 @@ use crate::app::storage_compat::ECStore;
use crate::app::storage_compat::StorageError;
use crate::app::storage_compat::get_global_notification_sys;
use crate::app::storage_compat::object_api_utils::to_s3s_etag;
use crate::app::storage_compat::{AppObjectLockConfigExt as _, AppVersioningConfigExt as _};
use crate::app::storage_compat::{
bucket_target_sys::BucketTargetSys,
lifecycle::bucket_lifecycle_ops::{
@@ -57,8 +58,6 @@ use futures::StreamExt;
use http::StatusCode;
use metrics::counter;
use rustfs_config::RUSTFS_REGION;
use rustfs_ecstore::api::bucket::object_lock::ObjectLockApi as _;
use rustfs_ecstore::api::bucket::versioning::VersioningApi as _;
use rustfs_madmin::{SITE_REPL_API_VERSION, SRBucketMeta};
use rustfs_policy::policy::{
action::{Action, S3Action},
+5 -6
View File
@@ -52,6 +52,9 @@ use crate::app::storage_compat::ECStore;
use crate::app::storage_compat::object_api_utils::to_s3s_etag;
use crate::app::storage_compat::quota::checker::QuotaChecker;
use crate::app::storage_compat::storageclass;
use crate::app::storage_compat::{
AppReplicationConfigExt as _, AppVersioningConfigExt as _, predict_lifecycle_expiration, validate_restore_request,
};
use crate::app::storage_compat::{DiskError, is_all_buckets_not_found};
use crate::app::storage_compat::{DynReader, HashReader, WritePlan, wrap_reader};
use crate::app::storage_compat::{
@@ -80,10 +83,6 @@ use crate::app::storage_compat::{
};
use crate::server::convert_ecstore_object_info;
use rustfs_concurrency::GetObjectQueueSnapshot;
use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::RestoreRequestOps as _;
use rustfs_ecstore::api::bucket::lifecycle::lifecycle::Lifecycle as _;
use rustfs_ecstore::api::bucket::replication::ReplicationConfigurationExt as _;
use rustfs_ecstore::api::bucket::versioning::VersioningApi as _;
use rustfs_filemeta::{
REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicateTargetDecision, ReplicationState, ReplicationStatusType,
ReplicationType, RestoreStatusOps, VersionPurgeStatusType, parse_restore_obj_status, replication_statuses_map,
@@ -1281,7 +1280,7 @@ async fn resolve_put_object_expiration(bucket: &str, obj_info: &ObjectInfo) -> O
};
let obj_opts = lifecycle::ObjectOpts::from_object_info(obj_info);
let event = lifecycle_config.predict_expiration(&obj_opts).await;
let event = predict_lifecycle_expiration(&lifecycle_config, &obj_opts).await;
debug!(
bucket,
action = ?event.action,
@@ -4305,7 +4304,7 @@ impl DefaultObjectUsecase {
}
// Validate restore request
if let Err(e) = rreq.validate(store.clone()) {
if let Err(e) = validate_restore_request(&rreq, store.clone()) {
return Err(S3Error::with_message(
S3ErrorCode::Custom("ErrValidRestoreObject".into()),
format!("Restore object validation failed: {}", e),
+53
View File
@@ -65,6 +65,59 @@ pub(crate) fn EndpointServerPools(pools: Vec<PoolEndpoints>) -> EndpointServerPo
crate::app::storage_compat::ecstore_layout::EndpointServerPools::from(pools)
}
pub(crate) trait AppObjectLockConfigExt {
fn enabled(&self) -> bool;
}
impl AppObjectLockConfigExt for s3s::dto::ObjectLockConfiguration {
fn enabled(&self) -> bool {
<s3s::dto::ObjectLockConfiguration as ecstore_bucket::object_lock::ObjectLockApi>::enabled(self)
}
}
pub(crate) trait AppReplicationConfigExt {
fn filter_target_arns(&self, obj: &replication::ObjectOpts) -> Vec<String>;
fn replicate(&self, opts: &replication::ObjectOpts) -> bool;
}
impl AppReplicationConfigExt for s3s::dto::ReplicationConfiguration {
fn filter_target_arns(&self, obj: &replication::ObjectOpts) -> Vec<String> {
<s3s::dto::ReplicationConfiguration as ecstore_bucket::replication::ReplicationConfigurationExt>::filter_target_arns(
self, obj,
)
}
fn replicate(&self, opts: &replication::ObjectOpts) -> bool {
<s3s::dto::ReplicationConfiguration as ecstore_bucket::replication::ReplicationConfigurationExt>::replicate(self, opts)
}
}
pub(crate) trait AppVersioningConfigExt {
fn prefix_enabled(&self, prefix: &str) -> bool;
fn suspended(&self) -> bool;
}
impl AppVersioningConfigExt for s3s::dto::VersioningConfiguration {
fn prefix_enabled(&self, prefix: &str) -> bool {
<s3s::dto::VersioningConfiguration as ecstore_bucket::versioning::VersioningApi>::prefix_enabled(self, prefix)
}
fn suspended(&self) -> bool {
<s3s::dto::VersioningConfiguration as ecstore_bucket::versioning::VersioningApi>::suspended(self)
}
}
pub(crate) async fn predict_lifecycle_expiration(
lifecycle: &s3s::dto::BucketLifecycleConfiguration,
obj: &lifecycle::lifecycle::ObjectOpts,
) -> lifecycle::lifecycle::Event {
ecstore_bucket::lifecycle::lifecycle::Lifecycle::predict_expiration(lifecycle, obj).await
}
pub(crate) fn validate_restore_request(request: &s3s::dto::RestoreRequest, api: Arc<ECStore>) -> std::io::Result<()> {
<s3s::dto::RestoreRequest as ecstore_bucket::lifecycle::bucket_lifecycle_ops::RestoreRequestOps>::validate(request, api)
}
pub(crate) async fn get_server_info(get_pools: bool) -> rustfs_madmin::InfoMessage {
crate::app::storage_compat::ecstore_admin::get_server_info(get_pools).await
}
+1 -2
View File
@@ -29,12 +29,11 @@ use crate::storage::storage_compat::{
is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, record_replication_proxy, serialize,
update_bucket_metadata_config,
};
use crate::storage::storage_compat::{StorageReplicationConfigExt as _, StorageVersioningConfigExt as _};
use crate::storage::{parse_object_lock_legal_hold, parse_object_lock_retention, validate_bucket_object_lock_enabled};
use crate::table_catalog;
use http::StatusCode;
use metrics::{counter, histogram};
use rustfs_ecstore::api::bucket::replication::ReplicationConfigurationExt as _;
use rustfs_ecstore::api::bucket::versioning::VersioningApi as _;
use rustfs_io_metrics::record_s3_op;
use rustfs_s3_ops::S3Operation;
use rustfs_storage_api::{BucketOperations, BucketOptions, ObjectLockRetentionOptions, ObjectOperations as _};
+1 -1
View File
@@ -16,6 +16,7 @@ use crate::config::{RustFSBufferConfig, WorkloadProfile, get_global_buffer_confi
use crate::error::ApiError;
use crate::server::cors;
use crate::storage::ecfs::ListObjectUnorderedQuery;
use crate::storage::storage_compat::StorageReplicationConfigExt as _;
use crate::storage::storage_compat::{
StorageError, add_object_lock_years, get_bucket_cors_config, get_bucket_object_lock_config, get_bucket_replication_config,
resolve_object_store_handle,
@@ -23,7 +24,6 @@ use crate::storage::storage_compat::{
use http::header::{IF_MATCH, IF_MODIFIED_SINCE, IF_NONE_MATCH, IF_UNMODIFIED_SINCE};
use http::{HeaderMap, HeaderValue, StatusCode};
use metrics::counter;
use rustfs_ecstore::api::bucket::replication::ReplicationConfigurationExt as _;
use rustfs_storage_api::{BucketOperations, BucketOptions};
use rustfs_targets::EventName;
use rustfs_targets::arn::{TargetID, TargetIDError};
+22
View File
@@ -242,6 +242,28 @@ pub(crate) async fn find_local_disk_by_ref(disk_ref: &str) -> Option<DiskStore>
ecstore_storage::find_local_disk_by_ref(disk_ref).await
}
pub(crate) trait StorageReplicationConfigExt {
fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool;
}
impl StorageReplicationConfigExt for s3s::dto::ReplicationConfiguration {
fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool {
<s3s::dto::ReplicationConfiguration as ecstore_bucket::replication::ReplicationConfigurationExt>::has_active_rules(
self, prefix, recursive,
)
}
}
pub(crate) trait StorageVersioningConfigExt {
fn enabled(&self) -> bool;
}
impl StorageVersioningConfigExt for s3s::dto::VersioningConfiguration {
fn enabled(&self) -> bool {
<s3s::dto::VersioningConfiguration as ecstore_bucket::versioning::VersioningApi>::enabled(self)
}
}
pub(crate) type GetObjectReader = <ECStore as rustfs_storage_api::ObjectIO>::GetObjectReader;
pub(crate) type ObjectInfo = <ECStore as rustfs_storage_api::ObjectOperations>::ObjectInfo;
pub(crate) type ObjectOptions = <ECStore as rustfs_storage_api::ObjectOperations>::ObjectOptions;