refactor: prune trait import compat re-exports (#3699)

This commit is contained in:
安正超
2026-06-21 23:17:58 +08:00
committed by GitHub
parent f35821f75d
commit 6b3d96fde3
29 changed files with 272 additions and 137 deletions
@@ -14,8 +14,8 @@
// limitations under the License.
use crate::common::workspace_root;
use crate::storage_compat::{TonicInterceptor, gen_tonic_signature_interceptor, node_service_time_out_client};
use crate::storage_compat::{VolumeInfo, WalkDirOptions};
use crate::storage_compat::{gen_tonic_signature_interceptor, node_service_time_out_client};
use futures::future::join_all;
use rmp_serde::{Deserializer, Serializer};
use rustfs_filemeta::{MetaCacheEntry, MetacacheReader, MetacacheWriter};
@@ -53,9 +53,7 @@ async fn ping() -> Result<(), Box<dyn Error>> {
assert!(decoded_payload.is_ok());
// Create client
let mut client =
node_service_time_out_client(&CLUSTER_ADDR.to_string(), TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await?;
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), gen_tonic_signature_interceptor()).await?;
// Construct PingRequest
let request = Request::new(PingRequest {
@@ -80,9 +78,7 @@ async fn ping() -> Result<(), Box<dyn Error>> {
#[tokio::test]
#[ignore = "requires running RustFS server at localhost:9000"]
async fn make_volume() -> Result<(), Box<dyn Error>> {
let mut client =
node_service_time_out_client(&CLUSTER_ADDR.to_string(), TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await?;
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), gen_tonic_signature_interceptor()).await?;
let request = Request::new(MakeVolumeRequest {
disk: "data".to_string(),
volume: "dandan".to_string(),
@@ -100,9 +96,7 @@ async fn make_volume() -> Result<(), Box<dyn Error>> {
#[tokio::test]
#[ignore = "requires running RustFS server at localhost:9000"]
async fn list_volumes() -> Result<(), Box<dyn Error>> {
let mut client =
node_service_time_out_client(&CLUSTER_ADDR.to_string(), TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await?;
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), gen_tonic_signature_interceptor()).await?;
let request = Request::new(ListVolumesRequest {
disk: "data".to_string(),
});
@@ -132,9 +126,7 @@ async fn walk_dir() -> Result<(), Box<dyn Error>> {
let (rd, mut wr) = tokio::io::duplex(1024);
let mut buf = Vec::new();
opts.serialize(&mut Serializer::new(&mut buf))?;
let mut client =
node_service_time_out_client(&CLUSTER_ADDR.to_string(), TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await?;
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), gen_tonic_signature_interceptor()).await?;
let disk_path = std::env::var_os("RUSTFS_DISK_PATH").map(PathBuf::from).unwrap_or_else(|| {
let mut path = workspace_root();
path.push("target");
@@ -187,9 +179,7 @@ async fn walk_dir() -> Result<(), Box<dyn Error>> {
#[tokio::test]
#[ignore = "requires running RustFS server at localhost:9000"]
async fn read_all() -> Result<(), Box<dyn Error>> {
let mut client =
node_service_time_out_client(&CLUSTER_ADDR.to_string(), TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await?;
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), gen_tonic_signature_interceptor()).await?;
let request = Request::new(ReadAllRequest {
disk: "data".to_string(),
volume: "ff".to_string(),
@@ -207,9 +197,7 @@ async fn read_all() -> Result<(), Box<dyn Error>> {
#[tokio::test]
#[ignore = "requires running RustFS server at localhost:9000"]
async fn storage_info() -> Result<(), Box<dyn Error>> {
let mut client =
node_service_time_out_client(&CLUSTER_ADDR.to_string(), TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await?;
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string(), gen_tonic_signature_interceptor()).await?;
let request = Request::new(LocalStorageInfoRequest { metrics: true });
let response = client.local_storage_info(request).await?.into_inner();
+3 -1
View File
@@ -19,7 +19,9 @@ pub(crate) type TonicInterceptor = rustfs_ecstore::api::rpc::TonicInterceptor;
pub(crate) type VolumeInfo = rustfs_ecstore::api::disk::VolumeInfo;
pub(crate) type WalkDirOptions = rustfs_ecstore::api::disk::WalkDirOptions;
pub(crate) use rustfs_ecstore::api::rpc::gen_tonic_signature_interceptor;
pub(crate) fn gen_tonic_signature_interceptor() -> TonicInterceptor {
TonicInterceptor::Signature(rustfs_ecstore::api::rpc::gen_tonic_signature_interceptor())
}
pub(crate) async fn node_service_time_out_client(
addr: &String,
+1 -1
View File
@@ -42,7 +42,7 @@ pub mod capacity {
}
pub mod client {
pub use crate::client::{admin_handler_utils, object_api_utils, transition_api};
pub use crate::client::{admin_handler_utils, api_put_object, object_api_utils, transition_api};
}
pub mod cluster {
+2 -1
View File
@@ -20,6 +20,7 @@ use crate::heal::{
use crate::{Error, Result};
use metrics::{counter, gauge};
use rustfs_common::heal_channel::{HealAdmissionDropReason, HealAdmissionResult, HealRequestSource};
use rustfs_ecstore::api::disk::DiskAPI as _;
use rustfs_madmin::heal_commands::HealResultItem;
use std::{
collections::{BinaryHeap, HashMap},
@@ -33,7 +34,7 @@ use tokio::{
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
use super::storage_compat::{DiskAPI, DiskError, GLOBAL_LOCAL_DISK_MAP};
use super::storage_compat::{DiskError, GLOBAL_LOCAL_DISK_MAP};
const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60);
const LOG_COMPONENT_HEAL: &str = "heal";
+2 -1
View File
@@ -13,6 +13,7 @@
// limitations under the License.
use crate::{Error, Result};
use rustfs_ecstore::api::disk::DiskAPI as _;
use serde::{Deserialize, Serialize};
use std::path::Path;
use std::sync::Arc;
@@ -21,7 +22,7 @@ use tokio::sync::RwLock;
use tracing::{debug, warn};
use uuid::Uuid;
use super::storage_compat::{BUCKET_META_PREFIX, DiskAPI, DiskError, DiskStore, RUSTFS_META_BUCKET};
use super::storage_compat::{BUCKET_META_PREFIX, DiskError, DiskStore, RUSTFS_META_BUCKET};
const LOG_COMPONENT_HEAL: &str = "heal";
const LOG_SUBSYSTEM_RESUME: &str = "resume";
+10 -2
View File
@@ -22,9 +22,17 @@ pub(crate) type ECStore = rustfs_ecstore::api::storage::ECStore;
pub(crate) type EcstoreError = rustfs_ecstore::api::error::Error;
pub(crate) type Endpoint = rustfs_ecstore::api::disk::endpoint::Endpoint;
pub(crate) type StorageError = rustfs_ecstore::api::error::StorageError;
pub(crate) type LocalDiskMap = std::collections::HashMap<String, Option<DiskStore>>;
pub(crate) use rustfs_ecstore::api::disk::DiskAPI;
pub(crate) use rustfs_ecstore::api::global::GLOBAL_LOCAL_DISK_MAP;
pub(crate) struct GlobalLocalDiskMap;
pub(crate) static GLOBAL_LOCAL_DISK_MAP: GlobalLocalDiskMap = GlobalLocalDiskMap;
impl GlobalLocalDiskMap {
pub(crate) async fn read(&self) -> tokio::sync::RwLockReadGuard<'static, LocalDiskMap> {
rustfs_ecstore::api::global::GLOBAL_LOCAL_DISK_MAP.read().await
}
}
#[cfg(test)]
pub(crate) type DiskOption = rustfs_ecstore::api::disk::DiskOption;
+5 -1
View File
@@ -19,10 +19,14 @@ use std::sync::Arc;
pub(crate) type DiskStore = rustfs_ecstore::api::disk::DiskStore;
pub(crate) type ECStore = rustfs_ecstore::api::storage::ECStore;
pub(crate) type Endpoint = rustfs_ecstore::api::disk::endpoint::Endpoint;
pub(crate) type EndpointServerPools = rustfs_ecstore::api::layout::EndpointServerPools;
pub(crate) type Endpoints = rustfs_ecstore::api::layout::Endpoints;
pub(crate) type PoolEndpoints = rustfs_ecstore::api::layout::PoolEndpoints;
pub(crate) use rustfs_ecstore::api::layout::EndpointServerPools;
#[allow(non_snake_case)]
pub(crate) fn EndpointServerPools(pools: Vec<PoolEndpoints>) -> EndpointServerPools {
rustfs_ecstore::api::layout::EndpointServerPools::from(pools)
}
pub(crate) async fn init_bucket_metadata_sys(api: Arc<ECStore>, buckets: Vec<String>) {
rustfs_ecstore::api::bucket::metadata_sys::init_bucket_metadata_sys(api, buckets).await;
+4 -2
View File
@@ -39,6 +39,8 @@ use rustfs_config::{
ENV_SCANNER_CYCLE_MAX_OBJECTS,
};
use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS};
use rustfs_ecstore::api::bucket::lifecycle::lifecycle::Lifecycle as _;
use rustfs_ecstore::api::bucket::replication::ReplicationConfigurationExt as _;
use rustfs_storage_api::{BucketOperations, BucketOptions, NamespaceLocking as _};
use serde::{Deserialize, Serialize};
use tokio::sync::mpsc;
@@ -47,8 +49,8 @@ use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, instrument, warn};
use crate::storage_compat::{
ECStore, EcstoreError, Lifecycle as _, RUSTFS_META_BUCKET, ReplicationConfigurationExt as _, get_lifecycle_config,
get_replication_config, is_erasure_sd, read_config, replace_bucket_usage_memory_from_info, save_config,
ECStore, EcstoreError, RUSTFS_META_BUCKET, get_lifecycle_config, get_replication_config, is_erasure_sd, read_config,
replace_bucket_usage_memory_from_info, save_config,
};
const LOG_COMPONENT_SCANNER: &str = "scanner";
+9 -9
View File
@@ -39,6 +39,10 @@ use rustfs_common::metrics::{
IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, UpdateCurrentPathFn,
current_path_updater, global_metrics,
};
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_ecstore::api::disk::DiskAPI as _;
use rustfs_filemeta::{
MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ReplicateObjectInfo, ReplicationStatusType, ReplicationType,
};
@@ -51,10 +55,10 @@ use tokio_util::sync::CancellationToken;
use tracing::{debug, error, warn};
use crate::storage_compat::{
BucketVersioningSys, Disk, DiskAPI as _, DiskError, DiskInfoOptions, Evaluator, Event, GLOBAL_ExpiryState, LcEventSrc,
Lifecycle, ListPathRawOptions, ObjectOpts, ReplicationConfig, ReplicationConfigurationExt as _, ReplicationQueueAdmission,
StorageError, VersioningApi, apply_expiry_rule, apply_transition_rule, is_erasure, is_reserved_or_invalid_bucket,
list_path_raw, path2_bucket_object, path2_bucket_object_with_base_path, queue_replication_heal_internal,
BucketVersioningSys, Disk, DiskError, DiskInfoOptions, Evaluator, Event, LcEventSrc, ListPathRawOptions, ObjectOpts,
ReplicationConfig, ReplicationQueueAdmission, StorageError, apply_expiry_rule, apply_transition_rule,
enqueue_global_newer_noncurrent, is_erasure, is_reserved_or_invalid_bucket, list_path_raw, path2_bucket_object,
path2_bucket_object_with_base_path, queue_replication_heal_internal,
};
use crate::{ScannerObjectInfo as ObjectInfo, ScannerObjectToDelete as ObjectToDelete};
@@ -896,11 +900,7 @@ impl ScannerItem {
let action = event.action;
let count = u64::try_from(to_delete_objs.len()).unwrap_or(u64::MAX);
let done_ilm = Metrics::time_ilm(action);
let queued = GLOBAL_ExpiryState
.write()
.await
.enqueue_by_newer_noncurrent(&self.bucket, to_delete_objs, event, &LcEventSrc::Scanner)
.await;
let queued = enqueue_global_newer_noncurrent(&self.bucket, to_delete_objs, event, &LcEventSrc::Scanner).await;
if record_scanner_ilm_action_if_queued(global_metrics(), action, count, queued) {
done_ilm(count)();
remaining_versions = remaining_versions.saturating_sub(noncurrent_accounting.len());
+9 -10
View File
@@ -26,6 +26,10 @@ use rustfs_common::heal_channel::HealScanMode;
use rustfs_common::metrics::{Metric, Metrics, emit_scan_bucket_drive_complete, emit_scan_bucket_drive_partial, global_metrics};
#[cfg(test)]
use rustfs_config::{ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS, ENV_SCANNER_MAX_CONCURRENT_SET_SCANS};
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_ecstore::api::disk::DiskAPI as _;
use rustfs_filemeta::FileMeta;
use rustfs_storage_api::{BucketInfo, BucketOperations, BucketOptions, DiskSetSelector, StorageAdminApi};
use rustfs_utils::path::path_join_buf;
@@ -44,10 +48,9 @@ use tracing::{debug, error, warn};
use crate::ScannerObjectInfo as ObjectInfo;
use crate::storage_compat::{
BucketTargetSys, BucketVersioningSys, Disk, DiskAPI, DiskError, ECStore, EcstoreError as Error, EcstoreResult as Result,
GLOBAL_ExpiryState, GLOBAL_TierConfigMgr, Lifecycle, ReplicationConfig, ReplicationConfigurationExt, STORAGE_FORMAT_FILE,
SetDisks, StorageError, VersioningApi as _, get_lifecycle_config, get_object_lock_config, get_replication_config,
resolve_scanner_object_store_handle, storageclass,
BucketTargetSys, BucketVersioningSys, Disk, DiskError, ECStore, EcstoreError as Error, EcstoreResult as Result,
ReplicationConfig, STORAGE_FORMAT_FILE, SetDisks, StorageError, enqueue_global_free_version, get_lifecycle_config,
get_object_lock_config, get_replication_config, list_global_tiers, resolve_scanner_object_store_handle, storageclass,
};
pub(crate) const SCANNER_SKIP_FILE_ERROR: &str = "skip file";
@@ -1349,10 +1352,7 @@ impl ScannerIODisk for Disk {
let mut size_summary = SizeSummary::default();
let tiers = {
let tier_config_mgr = GLOBAL_TierConfigMgr.read().await;
tier_config_mgr.list_tiers()
};
let tiers = list_global_tiers().await;
for tier in tiers.iter() {
size_summary.tier_stats.insert(tier.name.clone(), TierStats::default());
@@ -1374,9 +1374,8 @@ impl ScannerIODisk for Disk {
item.apply_actions(object_infos, lock_config, &mut size_summary).await;
if !free_version_infos.is_empty() {
let mut expiry_state = GLOBAL_ExpiryState.write().await;
for oi in free_version_infos {
expiry_state.enqueue_free_version(oi).await;
enqueue_global_free_version(oi).await;
}
}
+83 -20
View File
@@ -23,32 +23,26 @@ pub(crate) const STORAGE_FORMAT_FILE: &str = rustfs_ecstore::api::disk::STORAGE_
pub(crate) const TRANSITION_COMPLETE: &str = rustfs_ecstore::api::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
pub(crate) type Disk = rustfs_ecstore::api::disk::Disk;
#[cfg(test)]
pub(crate) type DiskStore = rustfs_ecstore::api::disk::DiskStore;
pub(crate) type DiskError = rustfs_ecstore::api::disk::error::DiskError;
pub(crate) type ECStore = rustfs_ecstore::api::storage::ECStore;
pub(crate) type EcstoreError = rustfs_ecstore::api::error::Error;
pub(crate) type EcstoreResult<T> = rustfs_ecstore::api::error::Result<T>;
pub(crate) type ListPathRawOptions = rustfs_ecstore::api::cache::ListPathRawOptions;
pub(crate) type BucketTargetSys = rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys;
pub(crate) type BucketVersioningSys = rustfs_ecstore::api::bucket::versioning_sys::BucketVersioningSys;
pub(crate) type DiskInfoOptions = rustfs_ecstore::api::disk::DiskInfoOptions;
pub(crate) type Evaluator = rustfs_ecstore::api::bucket::lifecycle::evaluator::Evaluator;
pub(crate) type Event = rustfs_ecstore::api::bucket::lifecycle::lifecycle::Event;
pub(crate) type LcEventSrc = rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc;
pub(crate) type ObjectOpts = rustfs_ecstore::api::bucket::lifecycle::lifecycle::ObjectOpts;
pub(crate) type ReplicationConfig = rustfs_ecstore::api::bucket::replication::ReplicationConfig;
pub(crate) type ReplicationHealQueueResult = rustfs_ecstore::api::bucket::replication::ReplicationHealQueueResult;
pub(crate) type ReplicationQueueAdmission = rustfs_ecstore::api::bucket::replication::ReplicationQueueAdmission;
pub(crate) type SetDisks = rustfs_ecstore::api::set_disk::SetDisks;
pub(crate) type StorageError = rustfs_ecstore::api::error::StorageError;
pub(crate) use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys;
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc;
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::{
GLOBAL_ExpiryState, apply_expiry_rule, apply_transition_rule,
};
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::evaluator::Evaluator;
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::lifecycle::{Event, Lifecycle, ObjectOpts};
pub(crate) use rustfs_ecstore::api::bucket::metadata_sys::{
get_lifecycle_config, get_object_lock_config, get_replication_config,
};
pub(crate) use rustfs_ecstore::api::bucket::replication::{
ReplicationConfig, ReplicationConfigurationExt, ReplicationQueueAdmission, queue_replication_heal_internal,
};
pub(crate) use rustfs_ecstore::api::bucket::versioning::VersioningApi;
pub(crate) use rustfs_ecstore::api::bucket::versioning_sys::BucketVersioningSys;
pub(crate) use rustfs_ecstore::api::disk::{DiskAPI, DiskInfoOptions};
pub(crate) use rustfs_ecstore::api::global::GLOBAL_TierConfigMgr;
pub type ScannerGetObjectReader = <ECStore as ObjectIO>::GetObjectReader;
pub type ScannerObjectInfo = <ECStore as rustfs_storage_api::ObjectOperations>::ObjectInfo;
pub type ScannerObjectOptions = <ECStore as rustfs_storage_api::ObjectOperations>::ObjectOptions;
@@ -61,10 +55,79 @@ pub(crate) mod storageclass {
}
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::config::init as init_ecstore_config_for_scanner_tests;
pub(crate) fn init_ecstore_config_for_scanner_tests() {
rustfs_ecstore::api::config::init();
}
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::disk::{DiskOption, endpoint::Endpoint, new_disk};
pub(crate) type DiskOption = rustfs_ecstore::api::disk::DiskOption;
#[cfg(test)]
pub(crate) type Endpoint = rustfs_ecstore::api::disk::endpoint::Endpoint;
#[cfg(test)]
pub(crate) async fn new_disk(ep: &Endpoint, opt: &DiskOption) -> rustfs_ecstore::api::disk::error::Result<DiskStore> {
rustfs_ecstore::api::disk::new_disk(ep, opt).await
}
pub(crate) async fn get_lifecycle_config(
bucket: &str,
) -> EcstoreResult<(s3s::dto::BucketLifecycleConfiguration, time::OffsetDateTime)> {
rustfs_ecstore::api::bucket::metadata_sys::get_lifecycle_config(bucket).await
}
pub(crate) async fn get_object_lock_config(
bucket: &str,
) -> EcstoreResult<(s3s::dto::ObjectLockConfiguration, time::OffsetDateTime)> {
rustfs_ecstore::api::bucket::metadata_sys::get_object_lock_config(bucket).await
}
pub(crate) async fn get_replication_config(
bucket: &str,
) -> EcstoreResult<(s3s::dto::ReplicationConfiguration, time::OffsetDateTime)> {
rustfs_ecstore::api::bucket::metadata_sys::get_replication_config(bucket).await
}
pub(crate) async fn apply_transition_rule(event: &Event, src: &LcEventSrc, oi: &ScannerObjectInfo) -> bool {
rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::apply_transition_rule(event, src, oi).await
}
pub(crate) async fn apply_expiry_rule(event: &Event, src: &LcEventSrc, oi: &ScannerObjectInfo) -> bool {
rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::apply_expiry_rule(event, src, oi).await
}
pub(crate) async fn list_global_tiers() -> Vec<rustfs_ecstore::api::tier::tier_config::TierConfig> {
rustfs_ecstore::api::global::GLOBAL_TierConfigMgr.read().await.list_tiers()
}
pub(crate) async fn enqueue_global_free_version(oi: ScannerObjectInfo) {
rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::GLOBAL_ExpiryState
.write()
.await
.enqueue_free_version(oi)
.await;
}
pub(crate) async fn enqueue_global_newer_noncurrent(
bucket: &str,
to_delete_objs: Vec<ObjectToDelete>,
event: Event,
src: &LcEventSrc,
) -> bool {
rustfs_ecstore::api::bucket::lifecycle::bucket_lifecycle_ops::GLOBAL_ExpiryState
.write()
.await
.enqueue_by_newer_noncurrent(bucket, to_delete_objs, event, src)
.await
}
pub(crate) async fn queue_replication_heal_internal(
bucket: &str,
oi: ScannerObjectInfo,
rcfg: ReplicationConfig,
retry_count: u32,
) -> ReplicationHealQueueResult {
rustfs_ecstore::api::bucket::replication::queue_replication_heal_internal(bucket, oi, rcfg, retry_count).await
}
pub(crate) fn resolve_scanner_object_store_handle() -> Option<Arc<ECStore>> {
rustfs_ecstore::api::global::resolve_object_store_handle()
+27 -5
View File
@@ -14,6 +14,7 @@
#![allow(dead_code, unused_imports)]
use std::collections::HashMap;
use std::sync::Arc;
use time::OffsetDateTime;
@@ -25,21 +26,42 @@ pub(crate) type BucketVersioningSys = rustfs_ecstore::api::bucket::versioning_sy
pub(crate) type DiskOption = rustfs_ecstore::api::disk::DiskOption;
pub(crate) type ECStore = rustfs_ecstore::api::storage::ECStore;
pub(crate) type Endpoint = rustfs_ecstore::api::disk::endpoint::Endpoint;
pub(crate) type EndpointServerPools = rustfs_ecstore::api::layout::EndpointServerPools;
pub(crate) type Endpoints = rustfs_ecstore::api::layout::Endpoints;
pub(crate) type PoolEndpoints = rustfs_ecstore::api::layout::PoolEndpoints;
pub(crate) type ReadCloser = rustfs_ecstore::api::client::transition_api::ReadCloser;
pub(crate) type ReaderImpl = rustfs_ecstore::api::client::transition_api::ReaderImpl;
pub(crate) type TierConfig = rustfs_ecstore::api::tier::tier_config::TierConfig;
pub(crate) type TierConfigMgr = rustfs_ecstore::api::tier::tier::TierConfigMgr;
pub(crate) type TierMinIO = rustfs_ecstore::api::tier::tier_config::TierMinIO;
pub(crate) type TierType = rustfs_ecstore::api::tier::tier_config::TierType;
pub(crate) type TransitionOptions = rustfs_ecstore::api::bucket::lifecycle::lifecycle::TransitionOptions;
pub(crate) type WarmBackendGetOpts = rustfs_ecstore::api::tier::warm_backend::WarmBackendGetOpts;
pub(crate) use rustfs_ecstore::api::disk::DiskAPI;
pub(crate) use rustfs_ecstore::api::global::GLOBAL_TierConfigMgr;
pub(crate) use rustfs_ecstore::api::layout::EndpointServerPools;
pub(crate) use rustfs_ecstore::api::tier::warm_backend::WarmBackend;
pub(crate) use rustfs_ecstore::api::tier::warm_backend::build_transition_put_options;
#[allow(non_snake_case)]
pub(crate) fn EndpointServerPools(pools: Vec<PoolEndpoints>) -> EndpointServerPools {
rustfs_ecstore::api::layout::EndpointServerPools::from(pools)
}
pub(crate) struct GlobalTierConfigMgrCompat;
#[allow(non_upper_case_globals)]
pub(crate) static GLOBAL_TierConfigMgr: GlobalTierConfigMgrCompat = GlobalTierConfigMgrCompat;
impl std::ops::Deref for GlobalTierConfigMgrCompat {
type Target = Arc<tokio::sync::RwLock<TierConfigMgr>>;
fn deref(&self) -> &Self::Target {
&rustfs_ecstore::api::global::GLOBAL_TierConfigMgr
}
}
pub(crate) fn build_transition_put_options(
storage_class: String,
metadata: HashMap<String, String>,
) -> rustfs_ecstore::api::client::api_put_object::PutObjectOptions {
rustfs_ecstore::api::tier::warm_backend::build_transition_put_options(storage_class, metadata)
}
pub(crate) async fn enqueue_transition_for_existing_objects(
api: Arc<ECStore>,
@@ -15,14 +15,16 @@
mod common;
use crate::common::storage_compat::{
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskAPI, DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints,
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints,
GLOBAL_TierConfigMgr, PoolEndpoints, ReadCloser, ReaderImpl, STORAGE_FORMAT_FILE, TierConfig, TierMinIO, TierType,
TransitionOptions, WarmBackend, WarmBackendGetOpts, build_transition_put_options, enqueue_transition_for_existing_objects,
TransitionOptions, WarmBackendGetOpts, build_transition_put_options, enqueue_transition_for_existing_objects,
get_bucket_metadata, init_background_expiry, init_bucket_metadata_sys, init_local_disks, new_disk,
path2_bucket_object_with_base_path, update_bucket_metadata,
};
use futures::FutureExt;
use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT;
use rustfs_ecstore::api::disk::DiskAPI as _;
use rustfs_ecstore::api::tier::warm_backend::WarmBackend;
use rustfs_filemeta::FileMeta;
use rustfs_scanner::scanner_folder::ScannerItem;
use rustfs_scanner::scanner_io::ScannerIODisk;