// 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. // #730: disk abstractions still carry staged health and direct-I/O migration paths. #![allow(dead_code)] pub mod disk_store; pub mod endpoint; pub mod error; pub mod error_conv; pub mod error_reduce; pub mod format; pub mod fs; pub mod health_state; pub mod local; pub mod os; pub const RUSTFS_META_BUCKET: &str = ".rustfs.sys"; pub const MIGRATING_META_BUCKET: &str = ".minio.sys"; pub const RUSTFS_META_MULTIPART_BUCKET: &str = ".rustfs.sys/multipart"; pub const RUSTFS_META_TMP_BUCKET: &str = ".rustfs.sys/tmp"; pub const RUSTFS_META_TMP_DELETED_BUCKET: &str = ".rustfs.sys/tmp/.trash"; pub const BUCKET_META_PREFIX: &str = "buckets"; pub const FORMAT_CONFIG_FILE: &str = "format.json"; /// Per-disk marker present while an erasure-set heal is rebuilding this disk. /// `LocalDisk::disk_info` reports `healing = true` while the file exists. pub const HEALING_MARKER_PATH: &str = "healing.bin"; pub const STORAGE_FORMAT_FILE: &str = "xl.meta"; pub const STORAGE_FORMAT_FILE_BACKUP: &str = "xl.meta.bkp"; pub const PART_TRANSACTION_NEW_META: &str = "new.meta"; pub const PART_TRANSACTION_OLD_META: &str = "old.meta"; pub const PART_TRANSACTION_ROLLBACK: &str = "rollback"; pub fn part_transaction_path(part_path: &str) -> String { match part_path.rsplit_once('/') { Some((parent, name)) => format!("{parent}/.{name}.rustfs-txn"), None => format!(".{part_path}.rustfs-txn"), } } use crate::cluster::rpc::RemoteDisk; use crate::cluster::rpc::build_internode_data_transport_from_env; use crate::disk::disk_store::LocalDiskWrapper; use crate::disk::health_state::RuntimeDriveHealthState; use crate::disk::local::ScanGuard; use bytes::Bytes; use endpoint::Endpoint; use error::DiskError; use error::{Error, Result}; use local::LocalDisk; use rustfs_filemeta::{FileInfo, ObjectPartInfo, RawFileInfo}; use rustfs_madmin::info_commands::DiskMetrics; use serde::{Deserialize, Serialize}; use std::{fmt::Debug, path::PathBuf, sync::Arc, time::Duration}; use time::OffsetDateTime; use tokio::io::{AsyncRead, AsyncWrite}; use uuid::Uuid; pub type DiskStore = Arc; pub type FileReader = Box; pub type FileWriter = Box; #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize, Deserialize)] pub struct SnapshotLeaseToken(Uuid); impl SnapshotLeaseToken { pub fn new() -> Self { Self(Uuid::new_v4()) } pub fn from_slice(bytes: &[u8]) -> Result { let uuid = Uuid::from_slice(bytes).map_err(|_| Error::other("invalid snapshot lease token"))?; if uuid.is_nil() { return Err(Error::other("invalid snapshot lease token")); } Ok(Self(uuid)) } pub fn as_bytes(&self) -> &[u8; 16] { self.0.as_bytes() } } impl Default for SnapshotLeaseToken { fn default() -> Self { Self::new() } } #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum DataDirDeleteStatus { Deleted, Deferred, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum PartTransactionAction { Commit, Rollback, } #[derive(Clone, Copy, Debug)] pub struct MmapCopyStageMetrics { pub(crate) path: &'static str, pub(crate) access_check_stage: &'static str, pub(crate) path_resolve_stage: &'static str, pub(crate) metadata_lookup_stage: &'static str, pub(crate) metadata_validate_stage: &'static str, pub(crate) blocking_wait_stage: &'static str, pub(crate) blocking_task_stage: &'static str, pub(crate) file_open_stage: &'static str, pub(crate) mmap_map_stage: &'static str, pub(crate) mmap_copy_stage: &'static str, pub(crate) direct_read_copy_stage: &'static str, } #[derive(Debug)] pub enum Disk { Local(Box), Remote(Box), } impl Disk { pub(crate) async fn set_disk_id_state(&self, id: Option) -> Result<()> { match self { Disk::Local(local_disk) => { local_disk.set_disk_id_state(id).await; Ok(()) } Disk::Remote(remote_disk) => remote_disk.set_disk_id(id).await, } } } #[async_trait::async_trait] impl DiskAPI for Disk { fn to_string(&self) -> String { match self { Disk::Local(local_disk) => local_disk.to_string(), Disk::Remote(remote_disk) => remote_disk.to_string(), } } async fn is_online(&self) -> bool { match self { Disk::Local(local_disk) => local_disk.is_online().await, Disk::Remote(remote_disk) => remote_disk.is_online().await, } } fn is_local(&self) -> bool { match self { Disk::Local(local_disk) => local_disk.is_local(), Disk::Remote(remote_disk) => remote_disk.is_local(), } } fn host_name(&self) -> String { match self { Disk::Local(local_disk) => local_disk.host_name(), Disk::Remote(remote_disk) => remote_disk.host_name(), } } fn endpoint(&self) -> Endpoint { match self { Disk::Local(local_disk) => local_disk.endpoint(), Disk::Remote(remote_disk) => remote_disk.endpoint(), } } async fn close(&self) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.close().await, Disk::Remote(remote_disk) => remote_disk.close().await, } } async fn get_disk_id(&self) -> Result> { match self { Disk::Local(local_disk) => local_disk.get_disk_id().await, Disk::Remote(remote_disk) => remote_disk.get_disk_id().await, } } async fn set_disk_id(&self, id: Option) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.set_disk_id(id).await, Disk::Remote(remote_disk) => remote_disk.set_disk_id(id).await, } } fn path(&self) -> PathBuf { match self { Disk::Local(local_disk) => local_disk.path(), Disk::Remote(remote_disk) => remote_disk.path(), } } fn get_disk_location(&self) -> DiskLocation { match self { Disk::Local(local_disk) => local_disk.get_disk_location(), Disk::Remote(remote_disk) => remote_disk.get_disk_location(), } } #[tracing::instrument(level = "trace", skip_all)] async fn make_volume(&self, volume: &str) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.make_volume(volume).await, Disk::Remote(remote_disk) => remote_disk.make_volume(volume).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn make_volumes(&self, volumes: Vec<&str>) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.make_volumes(volumes).await, Disk::Remote(remote_disk) => remote_disk.make_volumes(volumes).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn list_volumes(&self) -> Result> { match self { Disk::Local(local_disk) => local_disk.list_volumes().await, Disk::Remote(remote_disk) => remote_disk.list_volumes().await, } } async fn stat_volume(&self, volume: &str) -> Result { match self { Disk::Local(local_disk) => local_disk.stat_volume(volume).await, Disk::Remote(remote_disk) => remote_disk.stat_volume(volume).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.delete_volume(volume, force_delete).await, Disk::Remote(remote_disk) => remote_disk.delete_volume(volume, force_delete).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.walk_dir(opts, wr).await, Disk::Remote(remote_disk) => remote_disk.walk_dir(opts, wr).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn delete_version( &self, volume: &str, path: &str, fi: FileInfo, force_del_marker: bool, opts: DeleteOptions, ) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.delete_version(volume, path, fi, force_del_marker, opts).await, Disk::Remote(remote_disk) => remote_disk.delete_version(volume, path, fi, force_del_marker, opts).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn delete_versions(&self, volume: &str, versions: Vec, opts: DeleteOptions) -> Vec> { match self { Disk::Local(local_disk) => local_disk.delete_versions(volume, versions, opts).await, Disk::Remote(remote_disk) => remote_disk.delete_versions(volume, versions, opts).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.delete_paths(volume, paths).await, Disk::Remote(remote_disk) => remote_disk.delete_paths(volume, paths).await, } } async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> Result { match self { Disk::Local(local_disk) => local_disk.acquire_snapshot_lease(volume, path).await, Disk::Remote(remote_disk) => remote_disk.acquire_snapshot_lease(volume, path).await, } } async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.release_snapshot_lease(volume, path, token).await, Disk::Remote(remote_disk) => remote_disk.release_snapshot_lease(volume, path, token).await, } } async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result { match self { Disk::Local(local_disk) => local_disk.renew_snapshot_lease(volume, path, token).await, Disk::Remote(remote_disk) => remote_disk.renew_snapshot_lease(volume, path, token).await, } } async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result { match self { Disk::Local(local_disk) => local_disk.delete_data_dir(volume, path, opts).await, Disk::Remote(remote_disk) => remote_disk.delete_data_dir(volume, path, opts).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.write_metadata(_org_volume, volume, path, fi).await, Disk::Remote(remote_disk) => remote_disk.write_metadata(_org_volume, volume, path, fi).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.update_metadata(volume, path, fi, opts).await, Disk::Remote(remote_disk) => remote_disk.update_metadata(volume, path, fi, opts).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_version( &self, _org_volume: &str, volume: &str, path: &str, version_id: &str, opts: &ReadOptions, ) -> Result { match self { Disk::Local(local_disk) => local_disk.read_version(_org_volume, volume, path, version_id, opts).await, Disk::Remote(remote_disk) => remote_disk.read_version(_org_volume, volume, path, version_id, opts).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result { match self { Disk::Local(local_disk) => local_disk.read_xl(volume, path, read_data).await, Disk::Remote(remote_disk) => remote_disk.read_xl(volume, path, read_data).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn rename_data( &self, src_volume: &str, src_path: &str, fi: FileInfo, dst_volume: &str, dst_path: &str, ) -> Result { match self { Disk::Local(local_disk) => local_disk.rename_data(src_volume, src_path, fi, dst_volume, dst_path).await, Disk::Remote(remote_disk) => remote_disk.rename_data(src_volume, src_path, fi, dst_volume, dst_path).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn list_dir(&self, _origvolume: &str, volume: &str, dir_path: &str, count: i32) -> Result> { match self { Disk::Local(local_disk) => local_disk.list_dir(_origvolume, volume, dir_path, count).await, Disk::Remote(remote_disk) => remote_disk.list_dir(_origvolume, volume, dir_path, count).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_file(&self, volume: &str, path: &str) -> Result { match self { Disk::Local(local_disk) => local_disk.read_file(volume, path).await, Disk::Remote(remote_disk) => remote_disk.read_file(volume, path).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_file_stream(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result { match self { Disk::Local(local_disk) => local_disk.read_file_stream(volume, path, offset, length).await, Disk::Remote(remote_disk) => remote_disk.read_file_stream(volume, path, offset, length).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_file_mmap_copy(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result { match self { Disk::Local(local_disk) => local_disk.read_file_mmap_copy(volume, path, offset, length).await, Disk::Remote(remote_disk) => remote_disk.read_file_mmap_copy(volume, path, offset, length).await, } } async fn read_file_mmap_copy_with_metrics( &self, volume: &str, path: &str, offset: usize, length: usize, metrics: Option, ) -> Result { match self { Disk::Local(local_disk) => { local_disk .read_file_mmap_copy_with_metrics(volume, path, offset, length, metrics) .await } Disk::Remote(remote_disk) => remote_disk.read_file_mmap_copy(volume, path, offset, length).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn append_file(&self, volume: &str, path: &str) -> Result { match self { Disk::Local(local_disk) => local_disk.append_file(volume, path).await, Disk::Remote(remote_disk) => remote_disk.append_file(volume, path).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, _file_size: i64) -> Result { match self { Disk::Local(local_disk) => local_disk.create_file(_origvolume, volume, path, _file_size).await, Disk::Remote(remote_disk) => remote_disk.create_file(_origvolume, volume, path, _file_size).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.rename_file(src_volume, src_path, dst_volume, dst_path).await, Disk::Remote(remote_disk) => remote_disk.rename_file(src_volume, src_path, dst_volume, dst_path).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_parts(&self, bucket: &str, paths: &[String]) -> Result> { match self { Disk::Local(local_disk) => local_disk.read_parts(bucket, paths).await, Disk::Remote(remote_disk) => remote_disk.read_parts(bucket, paths).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Bytes) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.rename_part(src_volume, src_path, dst_volume, dst_path, meta).await, Disk::Remote(remote_disk) => { remote_disk .rename_part(src_volume, src_path, dst_volume, dst_path, meta) .await } } } async fn prepare_part_transaction( &self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Bytes, ) -> Result<()> { match self { Disk::Local(local_disk) => { local_disk .prepare_part_transaction(src_volume, src_path, dst_volume, dst_path, meta) .await } Disk::Remote(remote_disk) => { remote_disk .prepare_part_transaction(src_volume, src_path, dst_volume, dst_path, meta) .await } } } async fn settle_part_transaction(&self, volume: &str, path: &str, action: PartTransactionAction) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.settle_part_transaction(volume, path, action).await, Disk::Remote(remote_disk) => remote_disk.settle_part_transaction(volume, path, action).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.delete(volume, path, opt).await, Disk::Remote(remote_disk) => remote_disk.delete(volume, path, opt).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn verify_file(&self, volume: &str, path: &str, fi: &FileInfo) -> Result { match self { Disk::Local(local_disk) => local_disk.verify_file(volume, path, fi).await, Disk::Remote(remote_disk) => remote_disk.verify_file(volume, path, fi).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn check_parts(&self, volume: &str, path: &str, fi: &FileInfo) -> Result { match self { Disk::Local(local_disk) => local_disk.check_parts(volume, path, fi).await, Disk::Remote(remote_disk) => remote_disk.check_parts(volume, path, fi).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_multiple(&self, req: ReadMultipleReq) -> Result> { match self { Disk::Local(local_disk) => local_disk.read_multiple(req).await, Disk::Remote(remote_disk) => remote_disk.read_multiple(req).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn write_all(&self, volume: &str, path: &str, data: Bytes) -> Result<()> { match self { Disk::Local(local_disk) => local_disk.write_all(volume, path, data).await, Disk::Remote(remote_disk) => remote_disk.write_all(volume, path, data).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn read_all(&self, volume: &str, path: &str) -> Result { match self { Disk::Local(local_disk) => local_disk.read_all(volume, path).await, Disk::Remote(remote_disk) => remote_disk.read_all(volume, path).await, } } #[tracing::instrument(level = "trace", skip_all)] async fn disk_info(&self, opts: &DiskInfoOptions) -> Result { match self { Disk::Local(local_disk) => local_disk.disk_info(opts).await, Disk::Remote(remote_disk) => remote_disk.disk_info(opts).await, } } fn start_scan(&self) -> ScanGuard { match self { Disk::Local(local_disk) => local_disk.start_scan(), Disk::Remote(remote_disk) => remote_disk.start_scan(), } } async fn read_metadata(&self, volume: &str, path: &str) -> Result { match self { Disk::Local(local_disk) => local_disk.read_metadata(volume, path).await, Disk::Remote(remote_disk) => remote_disk.read_metadata(volume, path).await, } } } impl Disk { pub async fn ns_scanner_server_epoch(&self) -> Result> { match self { Disk::Local(_) => Ok(None), Disk::Remote(remote_disk) => remote_disk.ns_scanner_server_epoch().await, } } pub async fn open_ns_scanner_stream(&self, request: NsScannerOpenRequest) -> Result { match self { Disk::Remote(remote_disk) => remote_disk.open_ns_scanner_stream(request).await, Disk::Local(_) => Err(Error::other("namespace scanner stream requires a remote disk")), } } pub fn runtime_state(&self) -> RuntimeDriveHealthState { match self { Disk::Local(local_disk) => local_disk.runtime_state(), Disk::Remote(remote_disk) => remote_disk.runtime_state(), } } pub fn offline_duration_secs(&self) -> Option { match self { Disk::Local(local_disk) => local_disk.offline_duration_secs(), Disk::Remote(remote_disk) => remote_disk.offline_duration_secs(), } } pub fn last_capacity_snapshot(&self) -> Option<(u64, u64, u64, u64)> { match self { Disk::Local(local_disk) => local_disk.last_capacity_snapshot(), Disk::Remote(remote_disk) => remote_disk.last_capacity_snapshot(), } } #[cfg(test)] pub fn health_check_enabled_for_test(&self) -> bool { match self { Disk::Local(local_disk) => local_disk.health_check_enabled_for_test(), Disk::Remote(remote_disk) => remote_disk.health_check_enabled_for_test(), } } pub fn record_capacity_probe(&self, total: u64, used: u64, free: u64) { match self { Disk::Local(local_disk) => local_disk.record_capacity_probe(total, used, free), Disk::Remote(remote_disk) => remote_disk.record_capacity_probe(total, used, free), } } #[cfg(test)] pub fn force_runtime_state_for_test(&self, state: RuntimeDriveHealthState) { match self { Disk::Local(local_disk) => local_disk.force_runtime_state_for_test(state), Disk::Remote(remote_disk) => remote_disk.force_runtime_state_for_test(state), } } } #[derive(Debug)] pub struct NsScannerOpenRequest { pub request_id: Uuid, pub server_epoch: Uuid, pub session_id: Uuid, pub session_sequence: u64, pub next_cycle: u64, pub leader_epoch: u64, pub body: Vec, pub stall_timeout: Option, } impl Disk { /// Reset drive health so `connect_load_init_formats` retries are not blocked by a prior /// transient mark-faulty (same disk handles are reused across retries). pub fn reset_health_for_store_init_retry(&self) { match self { Disk::Local(local_disk) => local_disk.reset_health_for_store_init_retry(), Disk::Remote(remote_disk) => remote_disk.reset_health_for_store_init_retry(), } } /// Enable health monitoring on this disk. /// Called after startup format loading completes so that remote peers /// have time to come online before being marked as faulty. pub fn enable_health_check(&self) { match self { Disk::Local(local_disk) => local_disk.enable_health_check(), Disk::Remote(remote_disk) => remote_disk.enable_health_check(), } } } pub async fn new_disk(ep: &Endpoint, opt: &DiskOption) -> Result { if ep.is_local { let s = LocalDisk::new(ep, opt.cleanup).await?; Ok(Arc::new(Disk::Local(Box::new(LocalDiskWrapper::new(Arc::new(s), opt.health_check))))) } else { let data_transport = build_internode_data_transport_from_env(); let remote_disk = RemoteDisk::new(ep, opt, data_transport?).await?; Ok(Arc::new(Disk::Remote(Box::new(remote_disk)))) } } #[async_trait::async_trait] pub trait DiskAPI: Debug + Send + Sync + 'static { fn to_string(&self) -> String; async fn is_online(&self) -> bool; fn is_local(&self) -> bool; // LastConn fn host_name(&self) -> String; fn endpoint(&self) -> Endpoint; async fn close(&self) -> Result<()>; async fn get_disk_id(&self) -> Result>; async fn set_disk_id(&self, id: Option) -> Result<()>; fn path(&self) -> PathBuf; fn get_disk_location(&self) -> DiskLocation; // Healing // DiskInfo // NSScanner // Volume operations. async fn make_volume(&self, volume: &str) -> Result<()>; async fn make_volumes(&self, volume: Vec<&str>) -> Result<()>; async fn list_volumes(&self) -> Result>; async fn stat_volume(&self, volume: &str) -> Result; /// Delete a volume (bucket directory). When `force_delete` is false a /// non-empty volume is refused with `VolumeNotEmpty` (non-recursive); when /// true it is removed recursively. Callers on the heal/dangling path must /// pass false so a misclassified bucket that still holds data cannot be /// wiped (backlog#799 B1). async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()>; // Concurrent read/write pipeline w <- MetaCacheEntry async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()>; // Metadata operations async fn delete_version( &self, volume: &str, path: &str, fi: FileInfo, force_del_marker: bool, opts: DeleteOptions, ) -> Result<()>; async fn delete_versions(&self, volume: &str, versions: Vec, opts: DeleteOptions) -> Vec>; async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()>; async fn acquire_snapshot_lease(&self, _volume: &str, _path: &str) -> Result { Err(Error::other("snapshot leases are not supported by this disk")) } async fn release_snapshot_lease(&self, _volume: &str, _path: &str, _token: SnapshotLeaseToken) -> Result<()> { Err(Error::other("snapshot leases are not supported by this disk")) } async fn renew_snapshot_lease(&self, _volume: &str, _path: &str, _token: SnapshotLeaseToken) -> Result { Err(Error::other("snapshot leases are not supported by this disk")) } async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result { self.delete(volume, path, opts).await?; Ok(DataDirDeleteStatus::Deleted) } async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()>; async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()>; async fn read_version( &self, org_volume: &str, volume: &str, path: &str, version_id: &str, opts: &ReadOptions, ) -> Result; async fn batch_read_version(&self, req: BatchReadVersionReq) -> Result> { batch_read_version_one_by_one(self, req).await } async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result; async fn read_metadata(&self, volume: &str, path: &str) -> Result; async fn rename_data( &self, src_volume: &str, src_path: &str, file_info: FileInfo, dst_volume: &str, dst_path: &str, ) -> Result; // File operations. // Read every file and directory within the folder async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> Result>; async fn read_file(&self, volume: &str, path: &str) -> Result; async fn read_file_stream(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result; /// File read using mmap-then-copy on Unix or an efficient read on non-Unix. async fn read_file_mmap_copy(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result; async fn read_file_mmap_copy_with_metrics( &self, volume: &str, path: &str, offset: usize, length: usize, _metrics: Option, ) -> Result { self.read_file_mmap_copy(volume, path, offset, length).await } /// Historical name for `read_file_mmap_copy`. #[deprecated( since = "1.0.0-beta.8", note = "use read_file_mmap_copy; this path copies mmap data into owned Bytes" )] async fn read_file_zero_copy(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result { self.read_file_mmap_copy(volume, path, offset, length).await } async fn append_file(&self, volume: &str, path: &str) -> Result; async fn create_file(&self, origvolume: &str, volume: &str, path: &str, file_size: i64) -> Result; // ReadFileStream async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()>; async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Bytes) -> Result<()>; async fn prepare_part_transaction( &self, _src_volume: &str, _src_path: &str, _dst_volume: &str, _dst_path: &str, _meta: Bytes, ) -> Result<()> { Err(DiskError::MethodNotAllowed) } async fn settle_part_transaction(&self, _volume: &str, _path: &str, _action: PartTransactionAction) -> Result<()> { Err(DiskError::MethodNotAllowed) } async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()>; // VerifyFile async fn verify_file(&self, volume: &str, path: &str, fi: &FileInfo) -> Result; // CheckParts async fn check_parts(&self, volume: &str, path: &str, fi: &FileInfo) -> Result; // StatInfoFile async fn read_parts(&self, bucket: &str, paths: &[String]) -> Result>; async fn read_multiple(&self, req: ReadMultipleReq) -> Result>; // CleanAbandonedData async fn write_all(&self, volume: &str, path: &str, data: Bytes) -> Result<()>; async fn read_all(&self, volume: &str, path: &str) -> Result; async fn disk_info(&self, opts: &DiskInfoOptions) -> Result; fn start_scan(&self) -> ScanGuard; } pub async fn batch_read_version_one_by_one(disk: &D, req: BatchReadVersionReq) -> Result> where D: DiskAPI + ?Sized, { validate_batch_read_version_item_count(req.items.len())?; let mut responses = Vec::with_capacity(req.items.len()); for (index, item) in req.items.iter().enumerate() { let response = match disk .read_version(&item.org_volume, &item.volume, &item.path, &item.version_id, &req.opts) .await { Ok(file_info) => BatchReadVersionResp { index, path: item.path.clone(), version_id: item.version_id.clone(), success: true, file_info, error: String::new(), }, Err(err) => BatchReadVersionResp { index, path: item.path.clone(), version_id: item.version_id.clone(), success: false, file_info: FileInfo::default(), error: err.to_string(), }, }; responses.push(response); } Ok(responses) } #[derive(Debug, Default, Serialize, Deserialize)] pub struct CheckPartsResp { pub results: Vec, } #[derive(Debug, Serialize, Deserialize, Default)] pub struct UpdateMetadataOpts { pub no_persistence: bool, pub replace_user_metadata: bool, } pub struct DiskLocation { pub pool_idx: Option, pub set_idx: Option, pub disk_idx: Option, } impl DiskLocation { pub fn valid(&self) -> bool { self.pool_idx.is_some() && self.set_idx.is_some() && self.disk_idx.is_some() } } #[derive(Debug, Default, Serialize, Deserialize)] pub struct DiskInfoOptions { pub disk_id: String, pub metrics: bool, pub noop: bool, } #[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)] pub struct DiskInfo { pub total: u64, pub free: u64, pub used: u64, pub used_inodes: u64, pub free_inodes: u64, pub major: u64, pub minor: u64, pub nr_requests: u64, pub fs_type: String, pub root_disk: bool, pub healing: bool, pub scanning: bool, pub endpoint: String, pub mount_path: String, /// Leaf physical block devices backing this mount path when available. pub physical_device_ids: Vec, pub id: Option, pub rotational: bool, pub metrics: DiskMetrics, pub error: String, } #[derive(Clone, Debug, Default)] pub struct Info { pub total: u64, pub free: u64, pub used: u64, pub files: u64, pub ffree: u64, pub fstype: String, pub major: u64, pub minor: u64, pub name: String, pub rotational: bool, pub nrrequests: u64, } #[derive(Debug, Default, Clone, Serialize, Deserialize)] pub struct FileInfoVersions { // Name of the volume. pub volume: String, // Name of the file. pub name: String, // Represents the latest mod time of the // latest version. pub latest_mod_time: Option, pub versions: Vec, pub free_versions: Vec, } impl FileInfoVersions { pub fn find_version_index(&self, v: &str) -> Option { if v.is_empty() { return None; } let vid = Uuid::parse_str(v).unwrap_or_default(); self.versions.iter().position(|v| v.version_id == Some(vid)) } } #[derive(Debug, Default, Clone, Serialize, Deserialize)] pub struct WalkDirOptions { // Bucket to scanner pub bucket: String, // Directory inside the bucket. pub base_dir: String, // Do a full recursive scan. pub recursive: bool, // Include entries hidden by delete markers. #[serde(default)] pub incl_deleted: bool, // ReportNotFound will return errFileNotFound if all disks reports the BaseDir cannot be found. pub report_notfound: bool, // FilterPrefix will only return results with given prefix within folder. // Should never contain a slash. pub filter_prefix: Option, // ForwardTo will forward to the given object path. pub forward_to: Option, // Limit the number of returned objects if > 0. pub limit: i32, // DiskID contains the disk ID of the disk. // Leave empty to not check disk ID. pub disk_id: String, // Skip the wrapper-level total timeout for long streaming walks. #[serde(default)] pub skip_total_timeout: bool, // Override the wrapper-level total timeout for long background walks. #[serde(default)] pub timeout_ms: Option, // Override the remote stream stall timeout for long background walks. #[serde(default)] pub stall_timeout_ms: Option, } impl WalkDirOptions { pub fn timeout_duration(&self) -> Option { self.timeout_ms.map(std::time::Duration::from_millis) } pub fn stall_timeout_duration(&self) -> Option { self.stall_timeout_ms.map(std::time::Duration::from_millis) } } #[derive(Clone, Debug, Default)] pub struct DiskOption { pub cleanup: bool, pub health_check: bool, } /// Per-disk observation of the destination key's *current* (latest) version /// in the dst `xl.meta` that `rename_data` reads before it commits the /// incoming version (rustfs/backlog#1009). Mirrors the pre-PUT /// `get_object_info` outcome — which the app layer previously obtained with an /// extra full-disk fanout — bit for bit: `Absent` iff that lookup would report /// object-not-found (missing xl.meta, or only hidden free versions), and /// `Present(size)` iff it would return `Ok` with `ObjectInfo.size == size`. /// Note a delete-marker latest is `Present(0)`, not `Absent`: the lookup /// returns markers as `Ok(size 0)` and delete-marker accounting never /// decrements objects_count, so `Present(0)` is what keeps usage numbers /// identical. /// /// The parity target is the set-level (`SetDisks`) lookup. Two app-visible /// lookup variants deviated from it and from each other — the multi-pool /// `ECStore` path converts a delete-marker latest to not-found, and /// directory-key lookups pin the nil version — and both deviations /// over-incremented objects_count relative to delete-marker accounting, so /// the backfill standardizes on the self-consistent set-level answer. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] pub enum OldCurrentSize { /// The pre-PUT lookup would have reported object-not-found for this key. Absent, /// The pre-PUT lookup would have returned `Ok` with this `ObjectInfo.size` /// (0 for a delete-marker latest). Present(i64), } #[derive(Debug, Default, Serialize, Deserialize)] pub struct RenameDataResp { pub old_data_dir: Option, pub sign: Option>, /// `None` means unknown — the disk could not determine the previous /// current version (pre-#1009 peer on the wire, or an existing dst /// `xl.meta` that failed to parse). Consumers must treat unknown as /// "cannot vote", never as `Absent`. #[serde(default)] pub old_current_size: Option, } #[derive(Debug, Clone, Default, Serialize, Deserialize)] pub struct DeleteOptions { pub recursive: bool, pub immediate: bool, pub undo_write: bool, #[serde(default)] pub undo_delete: bool, pub old_data_dir: Option, } #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ReadMultipleReq { pub bucket: String, pub prefix: String, pub files: Vec, pub max_size: usize, pub metadata_only: bool, pub abort404: bool, pub max_results: usize, } #[derive(Debug, Clone, Default, Serialize, Deserialize)] pub struct ReadMultipleResp { pub bucket: String, pub prefix: String, pub file: String, pub exists: bool, pub error: String, pub data: Vec, pub mod_time: Option, } pub const BATCH_READ_VERSION_MAX_ITEMS: usize = 128; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct BatchReadVersionItem { pub org_volume: String, pub volume: String, pub path: String, pub version_id: String, } #[derive(Debug, Clone, Serialize, Deserialize)] pub struct BatchReadVersionReq { pub items: Vec, pub opts: ReadOptions, } #[derive(Debug, Clone, Default, Serialize, Deserialize)] pub struct BatchReadVersionResp { pub index: usize, pub path: String, pub version_id: String, pub success: bool, pub file_info: FileInfo, pub error: String, } pub fn validate_batch_read_version_item_count(item_count: usize) -> Result<()> { if item_count > BATCH_READ_VERSION_MAX_ITEMS { return Err(DiskError::other(format!( "batch read version item count {item_count} exceeds limit {BATCH_READ_VERSION_MAX_ITEMS}" ))); } Ok(()) } #[derive(Debug, Deserialize, Serialize)] pub struct VolumeInfo { pub name: String, pub created: Option, } #[derive(Deserialize, Serialize, Debug, Default, Clone)] pub struct ReadOptions { pub incl_free_versions: bool, pub read_data: bool, pub healing: bool, } pub const CHECK_PART_UNKNOWN: usize = 0; // Changing the order can cause a data loss // when running two nodes with incompatible versions pub const CHECK_PART_SUCCESS: usize = 1; pub const CHECK_PART_DISK_NOT_FOUND: usize = 2; pub const CHECK_PART_VOLUME_NOT_FOUND: usize = 3; pub const CHECK_PART_FILE_NOT_FOUND: usize = 4; pub const CHECK_PART_FILE_CORRUPT: usize = 5; pub fn conv_part_err_to_int(err: &Option) -> usize { match err { Some(DiskError::FileNotFound) | Some(DiskError::FileVersionNotFound) => CHECK_PART_FILE_NOT_FOUND, Some(DiskError::FileCorrupt) => CHECK_PART_FILE_CORRUPT, Some(DiskError::VolumeNotFound) => CHECK_PART_VOLUME_NOT_FOUND, Some(DiskError::DiskNotFound) => CHECK_PART_DISK_NOT_FOUND, None => CHECK_PART_SUCCESS, _ => { tracing::warn!("conv_part_err_to_int: unknown error: {err:?}"); CHECK_PART_UNKNOWN } } } pub fn has_part_err(part_errs: &[usize]) -> bool { part_errs.iter().any(|err| *err != CHECK_PART_SUCCESS) } pub fn count_part_not_success(part_errs: &[usize]) -> usize { part_errs.iter().filter(|err| **err != CHECK_PART_SUCCESS).count() } #[cfg(test)] mod tests { use super::*; use endpoint::Endpoint; use local::LocalDisk; use tokio::fs; use uuid::Uuid; /// Test DiskLocation validation #[test] fn test_disk_location_valid() { let valid_location = DiskLocation { pool_idx: Some(0), set_idx: Some(1), disk_idx: Some(2), }; assert!(valid_location.valid()); let invalid_location = DiskLocation { pool_idx: None, set_idx: None, disk_idx: None, }; assert!(!invalid_location.valid()); let partial_valid_location = DiskLocation { pool_idx: Some(0), set_idx: None, disk_idx: Some(2), }; assert!(!partial_valid_location.valid()); } /// Test FileInfoVersions find_version_index #[test] fn test_file_info_versions_find_version_index() { let mut versions = Vec::new(); let v1_uuid = Uuid::new_v4(); let v2_uuid = Uuid::new_v4(); let fi1 = FileInfo { version_id: Some(v1_uuid), ..Default::default() }; let fi2 = FileInfo { version_id: Some(v2_uuid), ..Default::default() }; versions.push(fi1); versions.push(fi2); let fiv = FileInfoVersions { volume: "test-bucket".to_string(), name: "test-object".to_string(), latest_mod_time: None, versions, free_versions: Vec::new(), }; assert_eq!(fiv.find_version_index(&v1_uuid.to_string()), Some(0)); assert_eq!(fiv.find_version_index(&v2_uuid.to_string()), Some(1)); assert_eq!(fiv.find_version_index("non-existent"), None); assert_eq!(fiv.find_version_index(""), None); } /// Test part error conversion functions #[test] fn test_conv_part_err_to_int() { assert_eq!(conv_part_err_to_int(&None), CHECK_PART_SUCCESS); assert_eq!( conv_part_err_to_int(&Some(Error::from(DiskError::DiskNotFound))), CHECK_PART_DISK_NOT_FOUND ); assert_eq!( conv_part_err_to_int(&Some(Error::from(DiskError::VolumeNotFound))), CHECK_PART_VOLUME_NOT_FOUND ); assert_eq!( conv_part_err_to_int(&Some(Error::from(DiskError::FileNotFound))), CHECK_PART_FILE_NOT_FOUND ); assert_eq!(conv_part_err_to_int(&Some(Error::from(DiskError::FileCorrupt))), CHECK_PART_FILE_CORRUPT); assert_eq!(conv_part_err_to_int(&Some(Error::from(DiskError::Unexpected))), CHECK_PART_UNKNOWN); } /// Test has_part_err function #[test] fn test_has_part_err() { assert!(!has_part_err(&[])); assert!(!has_part_err(&[CHECK_PART_SUCCESS])); assert!(!has_part_err(&[CHECK_PART_SUCCESS, CHECK_PART_SUCCESS])); assert!(has_part_err(&[CHECK_PART_FILE_NOT_FOUND])); assert!(has_part_err(&[CHECK_PART_SUCCESS, CHECK_PART_FILE_CORRUPT])); assert!(has_part_err(&[CHECK_PART_DISK_NOT_FOUND, CHECK_PART_VOLUME_NOT_FOUND])); } /// Test WalkDirOptions structure #[test] fn test_walk_dir_options() { let opts = WalkDirOptions { bucket: "test-bucket".to_string(), base_dir: "/path/to/dir".to_string(), recursive: true, incl_deleted: false, report_notfound: false, filter_prefix: Some("prefix_".to_string()), forward_to: Some("object/path".to_string()), limit: 100, disk_id: "disk-123".to_string(), skip_total_timeout: false, timeout_ms: Some(10_000), stall_timeout_ms: Some(20_000), }; assert_eq!(opts.bucket, "test-bucket"); assert_eq!(opts.base_dir, "/path/to/dir"); assert!(opts.recursive); assert!(!opts.incl_deleted); assert!(!opts.report_notfound); assert_eq!(opts.filter_prefix, Some("prefix_".to_string())); assert_eq!(opts.forward_to, Some("object/path".to_string())); assert_eq!(opts.limit, 100); assert_eq!(opts.disk_id, "disk-123"); assert!(!opts.skip_total_timeout); assert_eq!(opts.timeout_duration(), Some(std::time::Duration::from_secs(10))); assert_eq!(opts.stall_timeout_duration(), Some(std::time::Duration::from_secs(20))); } /// Test DeleteOptions structure #[test] fn test_delete_options() { let opts = DeleteOptions { recursive: true, immediate: false, undo_write: true, undo_delete: false, old_data_dir: Some(Uuid::new_v4()), }; assert!(opts.recursive); assert!(!opts.immediate); assert!(opts.undo_write); assert!(!opts.undo_delete); assert!(opts.old_data_dir.is_some()); } /// Test ReadOptions structure #[test] fn test_read_options() { let opts = ReadOptions { incl_free_versions: true, read_data: false, healing: true, }; assert!(opts.incl_free_versions); assert!(!opts.read_data); assert!(opts.healing); } /// Test UpdateMetadataOpts structure #[test] fn test_update_metadata_opts() { let opts = UpdateMetadataOpts { no_persistence: true, ..Default::default() }; assert!(opts.no_persistence); assert!(!opts.replace_user_metadata); } /// Test DiskOption structure #[test] fn test_disk_option() { let opt = DiskOption { cleanup: true, health_check: false, }; assert!(opt.cleanup); assert!(!opt.health_check); } /// Test DiskInfoOptions structure #[test] fn test_disk_info_options() { let opts = DiskInfoOptions { disk_id: "test-disk-id".to_string(), metrics: true, noop: false, }; assert_eq!(opts.disk_id, "test-disk-id"); assert!(opts.metrics); assert!(!opts.noop); } /// Test ReadMultipleReq structure #[test] fn test_read_multiple_req() { let req = ReadMultipleReq { bucket: "test-bucket".to_string(), prefix: "prefix/".to_string(), files: vec!["file1.txt".to_string(), "file2.txt".to_string()], max_size: 1024, metadata_only: false, abort404: true, max_results: 10, }; assert_eq!(req.bucket, "test-bucket"); assert_eq!(req.prefix, "prefix/"); assert_eq!(req.files.len(), 2); assert_eq!(req.max_size, 1024); assert!(!req.metadata_only); assert!(req.abort404); assert_eq!(req.max_results, 10); } /// Test ReadMultipleResp structure #[test] fn test_read_multiple_resp() { let resp = ReadMultipleResp { bucket: "test-bucket".to_string(), prefix: "prefix/".to_string(), file: "test-file.txt".to_string(), exists: true, error: "".to_string(), data: vec![1, 2, 3, 4], mod_time: Some(time::OffsetDateTime::now_utc()), }; assert_eq!(resp.bucket, "test-bucket"); assert_eq!(resp.prefix, "prefix/"); assert_eq!(resp.file, "test-file.txt"); assert!(resp.exists); assert!(resp.error.is_empty()); assert_eq!(resp.data, vec![1, 2, 3, 4]); assert!(resp.mod_time.is_some()); } /// Test VolumeInfo structure #[test] fn test_volume_info() { let now = time::OffsetDateTime::now_utc(); let vol_info = VolumeInfo { name: "test-volume".to_string(), created: Some(now), }; assert_eq!(vol_info.name, "test-volume"); assert_eq!(vol_info.created, Some(now)); } /// Test CheckPartsResp structure #[test] fn test_check_parts_resp() { let resp = CheckPartsResp { results: vec![CHECK_PART_SUCCESS, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_FILE_CORRUPT], }; assert_eq!(resp.results.len(), 3); assert_eq!(resp.results[0], CHECK_PART_SUCCESS); assert_eq!(resp.results[1], CHECK_PART_FILE_NOT_FOUND); assert_eq!(resp.results[2], CHECK_PART_FILE_CORRUPT); } /// Test RenameDataResp structure #[test] fn test_rename_data_resp() { let uuid = Uuid::new_v4(); let signature = vec![0x01, 0x02, 0x03]; let resp = RenameDataResp { old_data_dir: Some(uuid), sign: Some(signature.clone()), old_current_size: Some(OldCurrentSize::Present(42)), }; assert_eq!(resp.old_data_dir, Some(uuid)); assert_eq!(resp.sign, Some(signature)); assert_eq!(resp.old_current_size, Some(OldCurrentSize::Present(42))); } /// rustfs/backlog#1009: `old_current_size` must survive a named-msgpack /// round trip (the internode RPC encoding) for every variant. #[test] fn test_rename_data_resp_old_current_size_msgpack_roundtrip() { for old_current_size in [None, Some(OldCurrentSize::Absent), Some(OldCurrentSize::Present(1337))] { let resp = RenameDataResp { old_data_dir: Some(Uuid::new_v4()), sign: Some(vec![0x01, 0x02, 0x03]), old_current_size, }; let encoded = rmp_serde::encode::to_vec_named(&resp).expect("named msgpack should encode"); let decoded: RenameDataResp = rmp_serde::decode::from_slice(&encoded).expect("named msgpack should decode"); assert_eq!(decoded.old_data_dir, resp.old_data_dir); assert_eq!(decoded.sign, resp.sign); assert_eq!(decoded.old_current_size, resp.old_current_size); } } /// rustfs/backlog#1009: a payload from a peer that predates /// `old_current_size` must decode with the field defaulting to `None` /// (unknown), keeping mixed-version clusters wire-compatible. #[test] fn test_rename_data_resp_decodes_payload_without_old_current_size() { #[derive(Serialize)] struct LegacyRenameDataResp { old_data_dir: Option, sign: Option>, } let legacy = LegacyRenameDataResp { old_data_dir: Some(Uuid::new_v4()), sign: Some(vec![0x0a, 0x0b]), }; let encoded = rmp_serde::encode::to_vec_named(&legacy).expect("legacy named msgpack should encode"); let decoded: RenameDataResp = rmp_serde::decode::from_slice(&encoded).expect("legacy payload should decode"); assert_eq!(decoded.old_data_dir, legacy.old_data_dir); assert_eq!(decoded.sign, legacy.sign); assert_eq!(decoded.old_current_size, None); } /// Test constants #[test] fn test_constants() { assert_eq!(RUSTFS_META_BUCKET, ".rustfs.sys"); assert_eq!(RUSTFS_META_MULTIPART_BUCKET, ".rustfs.sys/multipart"); assert_eq!(RUSTFS_META_TMP_BUCKET, ".rustfs.sys/tmp"); assert_eq!(RUSTFS_META_TMP_DELETED_BUCKET, ".rustfs.sys/tmp/.trash"); assert_eq!(BUCKET_META_PREFIX, "buckets"); assert_eq!(FORMAT_CONFIG_FILE, "format.json"); assert_eq!(STORAGE_FORMAT_FILE, "xl.meta"); assert_eq!(STORAGE_FORMAT_FILE_BACKUP, "xl.meta.bkp"); assert_eq!(CHECK_PART_UNKNOWN, 0); assert_eq!(CHECK_PART_SUCCESS, 1); assert_eq!(CHECK_PART_DISK_NOT_FOUND, 2); assert_eq!(CHECK_PART_VOLUME_NOT_FOUND, 3); assert_eq!(CHECK_PART_FILE_NOT_FOUND, 4); assert_eq!(CHECK_PART_FILE_CORRUPT, 5); } /// Integration test for creating a local disk #[tokio::test] async fn test_new_disk_creation() { let test_dir = "./test_disk_creation"; fs::create_dir_all(&test_dir).await.unwrap(); let endpoint = Endpoint::try_from(test_dir).unwrap(); let opt = DiskOption { cleanup: false, health_check: true, }; let disk = new_disk(&endpoint, &opt).await; assert!(disk.is_ok()); let disk = disk.unwrap(); assert_eq!(disk.path(), rustfs_utils::canonicalize(test_dir).unwrap()); assert!(disk.is_local()); // Note: is_online() might return false for local disks without proper initialization // This is expected behavior for test environments // Clean up the test directory let _ = fs::remove_dir_all(&test_dir).await; } /// Test Disk enum pattern matching #[tokio::test] async fn test_disk_enum_methods() { let test_dir = "./test_disk_enum"; fs::create_dir_all(&test_dir).await.unwrap(); let endpoint = Endpoint::try_from(test_dir).unwrap(); let local_disk = LocalDisk::new(&endpoint, false).await.unwrap(); let disk = Disk::Local(Box::new(LocalDiskWrapper::new(Arc::new(local_disk), false))); // Test basic methods assert!(disk.is_local()); // Note: is_online() might return false for local disks without proper initialization // assert!(disk.is_online().await); // Note: host_name() for local disks might be empty or contain localhost // assert!(!disk.host_name().is_empty()); // Note: to_string() format might vary, so just check it's not empty assert!(!disk.to_string().is_empty()); // Test path method let path = disk.path(); assert!(path.exists()); // Test disk location let location = disk.get_disk_location(); assert!(location.valid() || (!location.valid() && endpoint.pool_idx < 0)); // Clean up the test directory let _ = fs::remove_dir_all(&test_dir).await; } #[tokio::test] #[serial_test::serial] async fn local_disk_id_state_does_not_publish_to_the_process_registry() { let local_dir = tempfile::tempdir().expect("local disk tempdir should be created"); let mut endpoint = Endpoint::try_from(local_dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); endpoint.set_pool_index(0); endpoint.set_set_index(0); endpoint.set_disk_index(0); let local_disk = LocalDisk::new(&endpoint, false).await.expect("local disk should initialize"); let disk = Disk::Local(Box::new(LocalDiskWrapper::new(Arc::new(local_disk), false))); let disk_id = Uuid::new_v4(); disk.set_disk_id_state(Some(disk_id)) .await .expect("local wrapper state should accept a disk ID"); let Disk::Local(local_disk) = &disk else { panic!("test disk should remain local"); }; assert_eq!(local_disk.get_current_disk_id().await, Some(disk_id)); assert!( !crate::runtime::global::current_ctx() .local_disk_id_map() .read() .await .contains_key(&disk_id), "state-only startup publication must not update the process disk-ID registry" ); disk.set_disk_id_state(None) .await .expect("local wrapper state should clear a disk ID"); assert_eq!(local_disk.get_current_disk_id().await, None); } #[tokio::test] async fn remote_disk_id_state_delegates_some_and_none() { let mut endpoint = Endpoint::try_from("http://remote-server:9000/data").expect("remote endpoint should parse"); endpoint.set_pool_index(0); endpoint.set_set_index(0); endpoint.set_disk_index(0); let remote_disk = RemoteDisk::new( &endpoint, &DiskOption { cleanup: false, health_check: false, }, Arc::new(crate::cluster::rpc::TcpHttpInternodeDataTransport), ) .await .expect("remote disk should initialize"); let disk = Disk::Remote(Box::new(remote_disk)); let disk_id = Uuid::new_v4(); disk.set_disk_id_state(Some(disk_id)) .await .expect("remote state should accept a disk ID"); assert_eq!(disk.get_disk_id().await.expect("remote disk ID should be readable"), Some(disk_id)); disk.set_disk_id_state(None) .await .expect("remote state should clear a disk ID"); assert_eq!(disk.get_disk_id().await.expect("remote disk ID should be readable"), None); } #[tokio::test] async fn reset_health_for_store_init_retry_delegates_to_disk_variants() { let local_dir = tempfile::tempdir().unwrap(); let local_path = local_dir.path().to_str().expect("tempdir path should be utf8"); let mut local_endpoint = Endpoint::try_from(local_path).expect("local endpoint should parse"); local_endpoint.set_pool_index(0); local_endpoint.set_set_index(0); local_endpoint.set_disk_index(0); let local_disk = LocalDisk::new(&local_endpoint, false).await.unwrap(); let local_disk = Disk::Local(Box::new(LocalDiskWrapper::new(Arc::new(local_disk), false))); let mut remote_endpoint = Endpoint::try_from("http://remote-server:9000/data").expect("remote endpoint should parse"); remote_endpoint.set_pool_index(0); remote_endpoint.set_set_index(0); remote_endpoint.set_disk_index(1); let remote_disk = RemoteDisk::new( &remote_endpoint, &DiskOption { cleanup: false, health_check: false, }, Arc::new(crate::cluster::rpc::TcpHttpInternodeDataTransport), ) .await .unwrap(); let remote_disk = Disk::Remote(Box::new(remote_disk)); for disk in [&local_disk, &remote_disk] { disk.force_runtime_state_for_test(RuntimeDriveHealthState::Offline); assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Offline); disk.reset_health_for_store_init_retry(); assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Online); } } }