mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-13 00:26:53 +00:00
refactor: wrap disk rpc compat trait methods (#3707)
This commit is contained in:
@@ -77,12 +77,14 @@ pub mod disk {
|
||||
pub use crate::disk::endpoint::Endpoint;
|
||||
pub use crate::disk::error::DiskError;
|
||||
pub use crate::disk::error_reduce::is_all_buckets_not_found;
|
||||
pub use crate::disk::local::ScanGuard;
|
||||
pub use crate::disk::{
|
||||
BUCKET_META_PREFIX, DeleteOptions, Disk, DiskAPI, DiskInfoOptions, DiskOption, DiskStore, FileInfoVersions,
|
||||
RUSTFS_META_BUCKET, ReadMultipleReq, ReadMultipleResp, ReadOptions, STORAGE_FORMAT_FILE, UpdateMetadataOpts, VolumeInfo,
|
||||
WalkDirOptions, new_disk,
|
||||
BUCKET_META_PREFIX, CheckPartsResp, DeleteOptions, Disk, DiskAPI, DiskInfo, DiskInfoOptions, DiskLocation, DiskOption,
|
||||
DiskStore, FileInfoVersions, FileReader, FileWriter, RUSTFS_META_BUCKET, ReadMultipleReq, ReadMultipleResp, ReadOptions,
|
||||
RenameDataResp, STORAGE_FORMAT_FILE, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, new_disk,
|
||||
};
|
||||
pub use crate::disk::{endpoint, error, error_reduce};
|
||||
pub use bytes::Bytes;
|
||||
}
|
||||
|
||||
pub mod error {
|
||||
|
||||
@@ -20,7 +20,6 @@ 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},
|
||||
@@ -34,7 +33,7 @@ use tokio::{
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
use super::storage_compat::{DiskError, GLOBAL_LOCAL_DISK_MAP};
|
||||
use super::storage_compat::{DiskError, GLOBAL_LOCAL_DISK_MAP, HealDiskExt as _};
|
||||
|
||||
const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60);
|
||||
const LOG_COMPONENT_HEAL: &str = "heal";
|
||||
|
||||
@@ -13,7 +13,6 @@
|
||||
// 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;
|
||||
@@ -22,7 +21,7 @@ use tokio::sync::RwLock;
|
||||
use tracing::{debug, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::storage_compat::{BUCKET_META_PREFIX, DiskError, DiskStore, RUSTFS_META_BUCKET};
|
||||
use super::storage_compat::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
|
||||
const LOG_COMPONENT_HEAL: &str = "heal";
|
||||
const LOG_SUBSYSTEM_RESUME: &str = "resume";
|
||||
|
||||
@@ -22,6 +22,7 @@ pub(crate) const BUCKET_META_PREFIX: &str = ecstore_disk::BUCKET_META_PREFIX;
|
||||
pub(crate) const RUSTFS_META_BUCKET: &str = ecstore_disk::RUSTFS_META_BUCKET;
|
||||
|
||||
pub(crate) type DiskError = ecstore_disk::error::DiskError;
|
||||
pub(crate) type DiskResult<T> = ecstore_disk::error::Result<T>;
|
||||
pub(crate) type DiskStore = ecstore_disk::DiskStore;
|
||||
pub(crate) type ECStore = ecstore_storage::ECStore;
|
||||
pub(crate) type EcstoreError = ecstore_error::Error;
|
||||
@@ -47,6 +48,51 @@ pub(crate) async fn new_disk(ep: &Endpoint, opt: &DiskOption) -> ecstore_disk::e
|
||||
ecstore_disk::new_disk(ep, opt).await
|
||||
}
|
||||
|
||||
pub(crate) trait HealDiskExt {
|
||||
fn endpoint(&self) -> Endpoint;
|
||||
async fn get_disk_id(&self) -> DiskResult<Option<uuid::Uuid>>;
|
||||
async fn read_all(&self, volume: &str, path: &str) -> DiskResult<ecstore_disk::Bytes>;
|
||||
async fn write_all(&self, volume: &str, path: &str, data: ecstore_disk::Bytes) -> DiskResult<()>;
|
||||
async fn delete(&self, volume: &str, path: &str, options: ecstore_disk::DeleteOptions) -> DiskResult<()>;
|
||||
async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> DiskResult<Vec<String>>;
|
||||
#[cfg(test)]
|
||||
async fn make_volume(&self, volume: &str) -> DiskResult<()>;
|
||||
}
|
||||
|
||||
impl<T> HealDiskExt for T
|
||||
where
|
||||
T: ecstore_disk::DiskAPI,
|
||||
{
|
||||
fn endpoint(&self) -> Endpoint {
|
||||
ecstore_disk::DiskAPI::endpoint(self)
|
||||
}
|
||||
|
||||
async fn get_disk_id(&self) -> DiskResult<Option<uuid::Uuid>> {
|
||||
ecstore_disk::DiskAPI::get_disk_id(self).await
|
||||
}
|
||||
|
||||
async fn read_all(&self, volume: &str, path: &str) -> DiskResult<ecstore_disk::Bytes> {
|
||||
ecstore_disk::DiskAPI::read_all(self, volume, path).await
|
||||
}
|
||||
|
||||
async fn write_all(&self, volume: &str, path: &str, data: ecstore_disk::Bytes) -> DiskResult<()> {
|
||||
ecstore_disk::DiskAPI::write_all(self, volume, path, data).await
|
||||
}
|
||||
|
||||
async fn delete(&self, volume: &str, path: &str, options: ecstore_disk::DeleteOptions) -> DiskResult<()> {
|
||||
ecstore_disk::DiskAPI::delete(self, volume, path, options).await
|
||||
}
|
||||
|
||||
async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> DiskResult<Vec<String>> {
|
||||
ecstore_disk::DiskAPI::list_dir(self, origvolume, volume, dir_path, count).await
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn make_volume(&self, volume: &str) -> DiskResult<()> {
|
||||
ecstore_disk::DiskAPI::make_volume(self, volume).await
|
||||
}
|
||||
}
|
||||
|
||||
pub type HealObjectInfo = <ECStore as rustfs_storage_api::ObjectOperations>::ObjectInfo;
|
||||
pub type HealObjectOptions = <ECStore as rustfs_storage_api::ObjectOperations>::ObjectOptions;
|
||||
pub type HealPutObjReader = <ECStore as rustfs_storage_api::ObjectIO>::PutObjectReader;
|
||||
|
||||
@@ -39,7 +39,6 @@ use rustfs_common::metrics::{
|
||||
IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource, UpdateCurrentPathFn,
|
||||
current_path_updater, global_metrics,
|
||||
};
|
||||
use rustfs_ecstore::api::disk::DiskAPI as _;
|
||||
use rustfs_filemeta::{
|
||||
MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ReplicateObjectInfo, ReplicationStatusType, ReplicationType,
|
||||
};
|
||||
@@ -53,10 +52,10 @@ use tracing::{debug, error, warn};
|
||||
|
||||
use crate::storage_compat::{
|
||||
BucketVersioningSys, Disk, DiskError, DiskInfoOptions, Evaluator, Event, LcEventSrc, ListPathRawOptions, ObjectOpts,
|
||||
ReplicationConfig, ReplicationQueueAdmission, ScannerLifecycleConfigExt as _, ScannerReplicationConfigExt as _,
|
||||
ScannerVersioningConfigExt as _, 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,
|
||||
ReplicationConfig, ReplicationQueueAdmission, ScannerDiskExt as _, ScannerLifecycleConfigExt as _,
|
||||
ScannerReplicationConfigExt as _, ScannerVersioningConfigExt as _, 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};
|
||||
|
||||
|
||||
@@ -26,7 +26,6 @@ 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::disk::DiskAPI as _;
|
||||
use rustfs_filemeta::FileMeta;
|
||||
use rustfs_storage_api::{BucketInfo, BucketOperations, BucketOptions, DiskSetSelector, StorageAdminApi};
|
||||
use rustfs_utils::path::path_join_buf;
|
||||
@@ -46,9 +45,10 @@ use tracing::{debug, error, warn};
|
||||
use crate::ScannerObjectInfo as ObjectInfo;
|
||||
use crate::storage_compat::{
|
||||
BucketTargetSys, BucketVersioningSys, Disk, DiskError, ECStore, EcstoreError as Error, EcstoreResult as Result,
|
||||
ReplicationConfig, STORAGE_FORMAT_FILE, ScannerLifecycleConfigExt as _, ScannerReplicationConfigExt as _,
|
||||
ScannerVersioningConfigExt as _, 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,
|
||||
ReplicationConfig, STORAGE_FORMAT_FILE, ScannerDiskExt as _, ScannerLifecycleConfigExt as _,
|
||||
ScannerReplicationConfigExt as _, ScannerVersioningConfigExt as _, 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";
|
||||
|
||||
@@ -19,6 +19,7 @@ use rustfs_ecstore::api::{
|
||||
set_disk as ecstore_set_disk, storage as ecstore_storage, tier as ecstore_tier,
|
||||
};
|
||||
use rustfs_storage_api::{HTTPRangeSpec, ObjectIO, ObjectToDelete};
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
@@ -30,7 +31,9 @@ pub(crate) const TRANSITION_COMPLETE: &str = ecstore_bucket::lifecycle::lifecycl
|
||||
pub(crate) type Disk = ecstore_disk::Disk;
|
||||
#[cfg(test)]
|
||||
pub(crate) type DiskStore = ecstore_disk::DiskStore;
|
||||
pub(crate) type DiskLocation = ecstore_disk::DiskLocation;
|
||||
pub(crate) type DiskError = ecstore_disk::error::DiskError;
|
||||
pub(crate) type DiskResult<T> = ecstore_disk::error::Result<T>;
|
||||
pub(crate) type ECStore = ecstore_storage::ECStore;
|
||||
pub(crate) type EcstoreError = ecstore_error::Error;
|
||||
pub(crate) type EcstoreResult<T> = ecstore_error::Result<T>;
|
||||
@@ -45,6 +48,7 @@ pub(crate) type ObjectOpts = ecstore_bucket::lifecycle::lifecycle::ObjectOpts;
|
||||
pub(crate) type ReplicationConfig = ecstore_bucket::replication::ReplicationConfig;
|
||||
pub(crate) type ReplicationHealQueueResult = ecstore_bucket::replication::ReplicationHealQueueResult;
|
||||
pub(crate) type ReplicationQueueAdmission = ecstore_bucket::replication::ReplicationQueueAdmission;
|
||||
pub(crate) type ScanGuard = ecstore_disk::ScanGuard;
|
||||
pub(crate) type SetDisks = ecstore_set_disk::SetDisks;
|
||||
pub(crate) type StorageError = ecstore_error::StorageError;
|
||||
|
||||
@@ -133,6 +137,39 @@ impl ScannerVersioningConfigExt for s3s::dto::VersioningConfiguration {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) trait ScannerDiskExt {
|
||||
async fn disk_info(&self, opts: &DiskInfoOptions) -> DiskResult<ecstore_disk::DiskInfo>;
|
||||
async fn read_metadata(&self, volume: &str, path: &str) -> DiskResult<ecstore_disk::Bytes>;
|
||||
fn path(&self) -> PathBuf;
|
||||
fn get_disk_location(&self) -> DiskLocation;
|
||||
fn start_scan(&self) -> ScanGuard;
|
||||
}
|
||||
|
||||
impl<T> ScannerDiskExt for T
|
||||
where
|
||||
T: ecstore_disk::DiskAPI,
|
||||
{
|
||||
async fn disk_info(&self, opts: &DiskInfoOptions) -> DiskResult<ecstore_disk::DiskInfo> {
|
||||
ecstore_disk::DiskAPI::disk_info(self, opts).await
|
||||
}
|
||||
|
||||
async fn read_metadata(&self, volume: &str, path: &str) -> DiskResult<ecstore_disk::Bytes> {
|
||||
ecstore_disk::DiskAPI::read_metadata(self, volume, path).await
|
||||
}
|
||||
|
||||
fn path(&self) -> PathBuf {
|
||||
ecstore_disk::DiskAPI::path(self)
|
||||
}
|
||||
|
||||
fn get_disk_location(&self) -> DiskLocation {
|
||||
ecstore_disk::DiskAPI::get_disk_location(self)
|
||||
}
|
||||
|
||||
fn start_scan(&self) -> ScanGuard {
|
||||
ecstore_disk::DiskAPI::start_scan(self)
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn apply_transition_rule(event: &Event, src: &LcEventSrc, oi: &ScannerObjectInfo) -> bool {
|
||||
ecstore_bucket::lifecycle::bucket_lifecycle_ops::apply_transition_rule(event, src, oi).await
|
||||
}
|
||||
|
||||
@@ -41,8 +41,22 @@ pub(crate) type TierConfigMgr = ecstore_tier::tier::TierConfigMgr;
|
||||
pub(crate) type TierMinIO = ecstore_tier::tier_config::TierMinIO;
|
||||
pub(crate) type TierType = ecstore_tier::tier_config::TierType;
|
||||
pub(crate) type TransitionOptions = ecstore_bucket::lifecycle::lifecycle::TransitionOptions;
|
||||
pub(crate) use ecstore_tier::warm_backend::WarmBackend as ScannerWarmBackend;
|
||||
pub(crate) type WarmBackendGetOpts = ecstore_tier::warm_backend::WarmBackendGetOpts;
|
||||
|
||||
pub(crate) trait ScannerTestDiskExt {
|
||||
async fn read_metadata(&self, volume: &str, path: &str) -> ecstore_disk::error::Result<ecstore_disk::Bytes>;
|
||||
}
|
||||
|
||||
impl<T> ScannerTestDiskExt for T
|
||||
where
|
||||
T: ecstore_disk::DiskAPI,
|
||||
{
|
||||
async fn read_metadata(&self, volume: &str, path: &str) -> ecstore_disk::error::Result<ecstore_disk::Bytes> {
|
||||
ecstore_disk::DiskAPI::read_metadata(self, volume, path).await
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(non_snake_case)]
|
||||
pub(crate) fn EndpointServerPools(pools: Vec<PoolEndpoints>) -> EndpointServerPools {
|
||||
ecstore_layout::EndpointServerPools::from(pools)
|
||||
|
||||
@@ -16,15 +16,13 @@ mod common;
|
||||
|
||||
use crate::common::storage_compat::{
|
||||
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints,
|
||||
GLOBAL_TierConfigMgr, PoolEndpoints, ReadCloser, ReaderImpl, STORAGE_FORMAT_FILE, TierConfig, TierMinIO, TierType,
|
||||
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,
|
||||
GLOBAL_TierConfigMgr, PoolEndpoints, ReadCloser, ReaderImpl, STORAGE_FORMAT_FILE, ScannerTestDiskExt as _,
|
||||
ScannerWarmBackend, TierConfig, TierMinIO, TierType, 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;
|
||||
@@ -687,7 +685,7 @@ impl MockWarmBackend {
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl WarmBackend for MockWarmBackend {
|
||||
impl ScannerWarmBackend for MockWarmBackend {
|
||||
async fn put(&self, object: &str, r: ReaderImpl, _length: i64) -> Result<String, std::io::Error> {
|
||||
let bytes = self.read_bytes(r).await?;
|
||||
Ok(self.put_bytes(object, bytes, HashMap::new()).await)
|
||||
|
||||
Reference in New Issue
Block a user