diff --git a/common/lock/src/local_locker.rs b/common/lock/src/local_locker.rs index 9802856a4..22ebfe0ba 100644 --- a/common/lock/src/local_locker.rs +++ b/common/lock/src/local_locker.rs @@ -275,14 +275,35 @@ impl Locker for LocalLocker { Ok(reply) } - async fn close(&self) {} + async fn refresh(&mut self, args: &LockArgs) -> Result { + let mut idx = 0; + let mut key = args.uid.to_string(); + format_uuid(&mut key, &idx); + match self.lock_uid.get(&key) { + Some(resource) => { + let mut resource = resource; + loop { + match self.lock_map.get_mut(resource) { + Some(_lris) => {} + None => { + let mut key = args.uid.to_string(); + format_uuid(&mut key, &0); + self.lock_uid.remove(&key); + return Ok(idx > 0); + } + } - async fn is_online(&self) -> bool { - true - } - - async fn is_local(&self) -> bool { - true + idx += 1; + let mut key = args.uid.to_string(); + format_uuid(&mut key, &idx); + resource = match self.lock_uid.get(&key) { + Some(resource) => resource, + None => return Ok(true), + }; + } + } + None => Ok(false), + } } // TODO: need add timeout mechanism @@ -350,37 +371,14 @@ impl Locker for LocalLocker { Ok(reply) } - async fn refresh(&mut self, args: &LockArgs) -> Result { - let mut idx = 0; - let mut key = args.uid.to_string(); - format_uuid(&mut key, &idx); - match self.lock_uid.get(&key) { - Some(resource) => { - let mut resource = resource; - loop { - match self.lock_map.get_mut(resource) { - Some(_lris) => {} - None => { - let mut key = args.uid.to_string(); - format_uuid(&mut key, &0); - self.lock_uid.remove(&key); - return Ok(idx > 0); - } - } + async fn close(&self) {} - idx += 1; - let mut key = args.uid.to_string(); - format_uuid(&mut key, &idx); - resource = match self.lock_uid.get(&key) { - Some(resource) => resource, - None => return Ok(true), - }; - } - } - None => { - return Ok(false); - } - } + async fn is_online(&self) -> bool { + true + } + + async fn is_local(&self) -> bool { + true } } diff --git a/common/lock/src/lrwmutex.rs b/common/lock/src/lrwmutex.rs index 9bc3415ef..79080e790 100644 --- a/common/lock/src/lrwmutex.rs +++ b/common/lock/src/lrwmutex.rs @@ -141,7 +141,7 @@ mod test { l_rw_lock.lock().await; - assert!(!(l_rw_lock.get_r_lock(id, source, &timeout).await)); + assert!(!l_rw_lock.get_r_lock(id, source, &timeout).await); l_rw_lock.un_lock().await; assert!(l_rw_lock.get_r_lock(id, source, &timeout).await); @@ -165,7 +165,7 @@ mod test { let two_fn = async { let two = Arc::clone(&l_rw_lock); let timeout = Duration::from_secs(2); - assert!(!(two.get_r_lock(id, source, &timeout).await)); + assert!(!two.get_r_lock(id, source, &timeout).await); sleep(Duration::from_secs(5)).await; assert!(two.get_r_lock(id, source, &timeout).await); two.un_r_lock().await; diff --git a/common/lock/src/remote_client.rs b/common/lock/src/remote_client.rs index eeafa96d0..3023cbc08 100644 --- a/common/lock/src/remote_client.rs +++ b/common/lock/src/remote_client.rs @@ -88,23 +88,6 @@ impl Locker for RemoteClient { Ok(response.success) } - async fn force_unlock(&mut self, args: &LockArgs) -> Result { - info!("remote force_unlock"); - let args = serde_json::to_string(args)?; - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(GenerallyLockRequest { args }); - - let response = client.force_un_lock(request).await?.into_inner(); - - if let Some(error_info) = response.error_info { - return Err(Error::from_string(error_info)); - } - - Ok(response.success) - } - async fn refresh(&mut self, args: &LockArgs) -> Result { info!("remote refresh"); let args = serde_json::to_string(args)?; @@ -122,8 +105,21 @@ impl Locker for RemoteClient { Ok(response.success) } - async fn is_local(&self) -> bool { - false + async fn force_unlock(&mut self, args: &LockArgs) -> Result { + info!("remote force_unlock"); + let args = serde_json::to_string(args)?; + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(GenerallyLockRequest { args }); + + let response = client.force_un_lock(request).await?.into_inner(); + + if let Some(error_info) = response.error_info { + return Err(Error::from_string(error_info)); + } + + Ok(response.success) } async fn close(&self) {} @@ -131,4 +127,8 @@ impl Locker for RemoteClient { async fn is_online(&self) -> bool { true } + + async fn is_local(&self) -> bool { + false + } } diff --git a/ecstore/src/admin_server_info.rs b/ecstore/src/admin_server_info.rs index 0c060888e..4d7e5168e 100644 --- a/ecstore/src/admin_server_info.rs +++ b/ecstore/src/admin_server_info.rs @@ -209,12 +209,12 @@ pub async fn get_server_info(get_pools: bool) -> InfoMessage { let mut versions = madmin::Versions::default(); let mut delete_markers = madmin::DeleteMarkers::default(); let mut usage = madmin::Usage::default(); - let mut mode = madmin::ITEM_INITIALIZING; + let mut mode = ITEM_INITIALIZING; let mut backend = madmin::ErasureBackend::default(); - let mut pools: HashMap> = HashMap::new(); + let mut pools: HashMap> = HashMap::new(); if let Some(store) = new_object_layer_fn() { - mode = madmin::ITEM_ONLINE; + mode = ITEM_ONLINE; match load_data_usage_from_backend(store.clone()).await { Ok(res) => { buckets.count = res.buckets_count; @@ -242,7 +242,7 @@ pub async fn get_server_info(get_pools: bool) -> InfoMessage { warn!("backend_info end {:?}", after4 - after3); - let mut all_disks: Vec = Vec::new(); + let mut all_disks: Vec = Vec::new(); for server in servers.iter() { all_disks.extend(server.disks.clone()); } diff --git a/ecstore/src/bitrot.rs b/ecstore/src/bitrot.rs index 05e55a430..c0b427e69 100644 --- a/ecstore/src/bitrot.rs +++ b/ecstore/src/bitrot.rs @@ -546,8 +546,7 @@ impl Writer for BitrotFileWriter { hasher.update(h_buf); hasher.finalize() }) - .await - .unwrap(); + .await?; if let Some(f) = self.inner.as_mut() { f.write_all(&hash_bytes).await?; @@ -775,7 +774,7 @@ mod test { if !algo.available() || *algo != BitrotAlgorithm::HighwayHash256 { continue; } - let checksum = decode_to_vec(checksums.get(algo).unwrap()).unwrap(); + let checksum = decode_to_vec(checksums.get(algo).unwrap())?; let mut h = algo.new_hasher(); let mut msg = Vec::with_capacity(h.size() * h.block_size()); diff --git a/ecstore/src/bucket/metadata.rs b/ecstore/src/bucket/metadata.rs index 135bbb785..a961072e4 100644 --- a/ecstore/src/bucket/metadata.rs +++ b/ecstore/src/bucket/metadata.rs @@ -288,12 +288,7 @@ impl BucketMetadata { } pub fn set_created(&mut self, created: Option) { - self.created = { - match created { - Some(t) => t, - None => OffsetDateTime::now_utc(), - } - } + self.created = { created.unwrap_or_else(|| OffsetDateTime::now_utc()) } } pub async fn save(&mut self) -> Result<()> { @@ -420,7 +415,6 @@ where #[cfg(test)] mod test { - use super::*; #[tokio::test] diff --git a/ecstore/src/bucket/metadata_sys.rs b/ecstore/src/bucket/metadata_sys.rs index 62dc0dcf9..73a1ff576 100644 --- a/ecstore/src/bucket/metadata_sys.rs +++ b/ecstore/src/bucket/metadata_sys.rs @@ -359,10 +359,10 @@ impl BucketMetadataSys { let bm = match load_bucket_metadata(self.api.clone(), bucket).await { Ok(res) => res, Err(err) => { - if *self.initialized.read().await { - return Err(Error::msg("errBucketMetadataNotInitialized")); + return if *self.initialized.read().await { + Err(Error::msg("errBucketMetadataNotInitialized")) } else { - return Err(err); + Err(err) } } }; @@ -381,11 +381,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_versioning_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Ok((VersioningConfiguration::default(), OffsetDateTime::UNIX_EPOCH)); + return if config::error::is_err_config_not_found(&err) { + Ok((VersioningConfiguration::default(), OffsetDateTime::UNIX_EPOCH)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -401,11 +401,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_bucket_policy err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketPolicyNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketPolicyNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -421,11 +421,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_tagging_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::TaggingNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::TaggingNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -441,11 +441,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_object_lock_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketObjectLockConfigNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketObjectLockConfigNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -461,11 +461,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_lifecycle_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketLifecycleNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketLifecycleNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -501,11 +501,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_sse_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketSSEConfigNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketSSEConfigNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -532,11 +532,11 @@ impl BucketMetadataSys { Ok((res, _)) => res, Err(err) => { warn!("get_quota_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketQuotaConfigNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketQuotaConfigNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -552,11 +552,11 @@ impl BucketMetadataSys { Ok(res) => res, Err(err) => { warn!("get_replication_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketReplicationConfigNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketReplicationConfigNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; @@ -576,11 +576,11 @@ impl BucketMetadataSys { Ok(res) => res, Err(err) => { warn!("get_replication_config err {:?}", &err); - if config::error::is_err_config_not_found(&err) { - return Err(Error::new(BucketMetadataError::BucketRemoteTargetNotFound)); + return if config::error::is_err_config_not_found(&err) { + Err(Error::new(BucketMetadataError::BucketRemoteTargetNotFound)) } else { - return Err(err); - } + Err(err) + }; } }; diff --git a/ecstore/src/bucket/versioning/mod.rs b/ecstore/src/bucket/versioning/mod.rs index fb3387cc9..1c0344f99 100644 --- a/ecstore/src/bucket/versioning/mod.rs +++ b/ecstore/src/bucket/versioning/mod.rs @@ -14,10 +14,6 @@ impl VersioningApi for VersioningConfiguration { fn enabled(&self) -> bool { self.status == Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)) } - fn suspended(&self) -> bool { - self.status == Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)) - } - fn prefix_enabled(&self, prefix: &str) -> bool { if self.status != Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)) { return false; @@ -46,6 +42,7 @@ impl VersioningApi for VersioningConfiguration { true } + fn prefix_suspended(&self, prefix: &str) -> bool { if self.status == Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)) { return true; @@ -79,4 +76,7 @@ impl VersioningApi for VersioningConfiguration { fn versioned(&self, prefix: &str) -> bool { self.prefix_enabled(prefix) || self.prefix_suspended(prefix) } + fn suspended(&self) -> bool { + self.status == Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)) + } } diff --git a/ecstore/src/config/com.rs b/ecstore/src/config/com.rs index 43e48c6d2..03038066e 100644 --- a/ecstore/src/config/com.rs +++ b/ecstore/src/config/com.rs @@ -119,14 +119,14 @@ pub async fn read_config_without_migrate(api: Arc) -> Result res, Err(err) => { - if is_err_config_not_found(&err) { + return if is_err_config_not_found(&err) { warn!("config not found, start to init"); let cfg = new_and_save_server_config(api).await?; warn!("config init done"); - return Ok(cfg); + Ok(cfg) } else { error!("read config err {:?}", &err); - return Err(err); + Err(err) } } }; @@ -141,14 +141,14 @@ async fn read_server_config(api: Arc, data: &[u8]) -> Result res, Err(err) => { - if is_err_config_not_found(&err) { + return if is_err_config_not_found(&err) { warn!("config not found init start"); let cfg = new_and_save_server_config(api).await?; warn!("config not found init done"); - return Ok(cfg); + Ok(cfg) } else { error!("read config err {:?}", &err); - return Err(err); + Err(err) } } }; @@ -189,10 +189,9 @@ async fn apply_dynamic_config(cfg: &mut Config, api: Arc) -> R async fn apply_dynamic_config_for_sub_sys(cfg: &mut Config, api: Arc, subsys: &str) -> Result<()> { let set_drive_counts = api.set_drive_counts(); if subsys == STORAGE_CLASS_SUB_SYS { - let kvs = match cfg.get_value(STORAGE_CLASS_SUB_SYS, DEFAULT_KV_KEY) { - Some(res) => res, - None => KVS::new(), - }; + let kvs = cfg + .get_value(STORAGE_CLASS_SUB_SYS, DEFAULT_KV_KEY) + .unwrap_or_else(|| KVS::new()); for (i, count) in set_drive_counts.iter().enumerate() { match storageclass::lookup_config(&kvs, *count) { diff --git a/ecstore/src/disk/endpoint.rs b/ecstore/src/disk/endpoint.rs index d901969c2..207e508f5 100644 --- a/ecstore/src/disk/endpoint.rs +++ b/ecstore/src/disk/endpoint.rs @@ -17,7 +17,7 @@ pub enum EndpointType { /// any type of endpoint. #[derive(Debug, PartialEq, Eq, Clone, Hash)] pub struct Endpoint { - pub url: url::Url, + pub url: Url, pub is_local: bool, pub pool_idx: i32, @@ -179,8 +179,8 @@ impl Endpoint { } } -/// parse a file path into an URL. -fn url_parse_from_file_path(value: &str) -> Result { +/// parse a file path into a URL. +fn url_parse_from_file_path(value: &str) -> Result { // Only check if the arg is an ip address and ask for scheme since its absent. // localhost, example.com, any FQDN cannot be disambiguated from a regular file path such as // /mnt/export1. So we go ahead and start the rustfs server in FS modes in these cases. @@ -202,7 +202,6 @@ fn url_parse_from_file_path(value: &str) -> Result { #[cfg(test)] mod test { - use super::*; #[test] @@ -215,10 +214,10 @@ mod test { expected_err: Option, } - let u2 = url::Url::parse("https://example.org/path").unwrap(); - let u4 = url::Url::parse("http://192.168.253.200/path").unwrap(); - let u6 = url::Url::parse("http://server:/path").unwrap(); - let root_slash_foo = url::Url::from_file_path("/foo").unwrap(); + let u2 = Url::parse("https://example.org/path").unwrap(); + let u4 = Url::parse("http://192.168.253.200/path").unwrap(); + let u6 = Url::parse("http://server:/path").unwrap(); + let root_slash_foo = Url::from_file_path("/foo").unwrap(); let test_cases = [ TestCase { diff --git a/ecstore/src/disk/error.rs b/ecstore/src/disk/error.rs index febf67d77..3c49980ac 100644 --- a/ecstore/src/disk/error.rs +++ b/ecstore/src/disk/error.rs @@ -301,8 +301,8 @@ pub fn clone_disk_err(e: &DiskError) -> Error { pub fn os_err_to_file_err(e: io::Error) -> Error { match e.kind() { - io::ErrorKind::NotFound => Error::new(DiskError::FileNotFound), - io::ErrorKind::PermissionDenied => Error::new(DiskError::FileAccessDenied), + ErrorKind::NotFound => Error::new(DiskError::FileNotFound), + ErrorKind::PermissionDenied => Error::new(DiskError::FileAccessDenied), // io::ErrorKind::ConnectionRefused => todo!(), // io::ErrorKind::ConnectionReset => todo!(), // io::ErrorKind::HostUnreachable => todo!(), @@ -350,7 +350,7 @@ pub fn os_err_to_file_err(e: io::Error) -> Error { pub struct FileAccessDeniedWithContext { pub path: PathBuf, #[source] - pub source: std::io::Error, + pub source: io::Error, } impl std::fmt::Display for FileAccessDeniedWithContext { diff --git a/ecstore/src/disk/format.rs b/ecstore/src/disk/format.rs index 0a3be1b28..05dbc7058 100644 --- a/ecstore/src/disk/format.rs +++ b/ecstore/src/disk/format.rs @@ -40,7 +40,7 @@ pub struct FormatErasureV3 { pub this: Uuid, /// Sets field carries the input disk order generated the first - /// time when fresh disks were supplied, it is a two dimensional + /// time when fresh disks were supplied, it is a two-dimensional /// array second dimension represents list of disks used per set. pub sets: Vec>, diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 8cbef7778..d07576f5c 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -57,6 +57,14 @@ impl DiskAPI for Disk { } } + #[tracing::instrument(skip(self))] + 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, + } + } + #[tracing::instrument(skip(self))] fn is_local(&self) -> bool { match self { @@ -73,14 +81,6 @@ impl DiskAPI for Disk { } } - #[tracing::instrument(skip(self))] - 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, - } - } - #[tracing::instrument(skip(self))] fn endpoint(&self) -> Endpoint { match self { @@ -97,22 +97,6 @@ impl DiskAPI for Disk { } } - #[tracing::instrument(skip(self))] - fn path(&self) -> PathBuf { - match self { - Disk::Local(local_disk) => local_disk.path(), - Disk::Remote(remote_disk) => remote_disk.path(), - } - } - - #[tracing::instrument(skip(self))] - 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(skip(self))] async fn get_disk_id(&self) -> Result> { match self { @@ -130,133 +114,18 @@ impl DiskAPI for Disk { } #[tracing::instrument(skip(self))] - async fn read_all(&self, volume: &str, path: &str) -> Result> { + fn path(&self) -> PathBuf { match self { - Disk::Local(local_disk) => local_disk.read_all(volume, path).await, - Disk::Remote(remote_disk) => remote_disk.read_all(volume, path).await, + Disk::Local(local_disk) => local_disk.path(), + Disk::Remote(remote_disk) => remote_disk.path(), } } #[tracing::instrument(skip(self))] - async fn write_all(&self, volume: &str, path: &str, data: Vec) -> Result<()> { + fn get_disk_location(&self) -> DiskLocation { 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(skip(self))] - 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(skip(self))] - 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(skip(self))] - 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(skip(self))] - async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec) -> 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 - } - } - } - - #[tracing::instrument(skip(self))] - 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(skip(self))] - async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, _file_size: usize) -> 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(skip(self))] - 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(skip(self))] - 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(skip(self))] - 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(skip(self))] - 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(skip(self, wr))] - 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(skip(self, fi))] - 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(skip(self))] - 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, + Disk::Local(local_disk) => local_disk.get_disk_location(), + Disk::Remote(remote_disk) => remote_disk.get_disk_location(), } } @@ -268,6 +137,14 @@ impl DiskAPI for Disk { } } + #[tracing::instrument(skip(self))] + 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(skip(self))] async fn list_volumes(&self) -> Result> { match self { @@ -285,49 +162,18 @@ impl DiskAPI for Disk { } #[tracing::instrument(skip(self))] - async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()> { + async fn delete_volume(&self, volume: &str) -> 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, + Disk::Local(local_disk) => local_disk.delete_volume(volume).await, + Disk::Remote(remote_disk) => remote_disk.delete_volume(volume).await, } } - #[tracing::instrument(skip(self))] - async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> { + #[tracing::instrument(skip(self, wr))] + async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> 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(skip(self))] - 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 = "debug", skip(self))] - 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(skip(self))] - 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, + Disk::Local(local_disk) => local_disk.walk_dir(opts, wr).await, + Disk::Remote(remote_disk) => remote_disk.walk_dir(opts, wr).await, } } @@ -359,6 +205,152 @@ impl DiskAPI for Disk { } } + #[tracing::instrument(skip(self))] + 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, + } + } + + #[tracing::instrument(skip(self))] + 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(skip(self))] + 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 = "debug", skip(self))] + 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(skip(self))] + 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(skip(self, fi))] + 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(skip(self))] + 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(skip(self))] + 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(skip(self))] + 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(skip(self))] + 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(skip(self))] + async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, _file_size: usize) -> 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(skip(self))] + 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(skip(self))] + async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec) -> 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 + } + } + } + + #[tracing::instrument(skip(self))] + 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(skip(self))] + 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(skip(self))] + 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(skip(self))] async fn read_multiple(&self, req: ReadMultipleReq) -> Result> { match self { @@ -368,10 +360,18 @@ impl DiskAPI for Disk { } #[tracing::instrument(skip(self))] - async fn delete_volume(&self, volume: &str) -> Result<()> { + async fn write_all(&self, volume: &str, path: &str, data: Vec) -> Result<()> { match self { - Disk::Local(local_disk) => local_disk.delete_volume(volume).await, - Disk::Remote(remote_disk) => remote_disk.delete_volume(volume).await, + 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(skip(self))] + 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, } } @@ -406,12 +406,12 @@ impl DiskAPI for Disk { } } -pub async fn new_disk(ep: &endpoint::Endpoint, opt: &DiskOption) -> Result { +pub async fn new_disk(ep: &Endpoint, opt: &DiskOption) -> Result { if ep.is_local { - let s = local::LocalDisk::new(ep, opt.cleanup).await?; + let s = LocalDisk::new(ep, opt.cleanup).await?; Ok(Arc::new(Disk::Local(Box::new(s)))) } else { - let remote_disk = remote::RemoteDisk::new(ep, opt).await?; + let remote_disk = RemoteDisk::new(ep, opt).await?; Ok(Arc::new(Disk::Remote(Box::new(remote_disk)))) } } @@ -763,7 +763,7 @@ impl MetaCacheEntry { let fi = fm.into_fileinfo(bucket, self.name.as_str(), "", false, false)?; - return Ok(fi); + Ok(fi) } pub fn file_info_versions(&self, bucket: &str) -> Result { diff --git a/ecstore/src/disk/remote.rs b/ecstore/src/disk/remote.rs index 2b4e64dc2..8bf552dc6 100644 --- a/ecstore/src/disk/remote.rs +++ b/ecstore/src/disk/remote.rs @@ -76,21 +76,21 @@ impl DiskAPI for RemoteDisk { } #[tracing::instrument(skip(self))] - fn is_local(&self) -> bool { + async fn is_online(&self) -> bool { + // TODO: 连接状态 + if node_service_time_out_client(&self.addr).await.is_ok() { + return true; + } false } #[tracing::instrument(skip(self))] - fn host_name(&self) -> String { - self.endpoint.host_port() + fn is_local(&self) -> bool { + false } #[tracing::instrument(skip(self))] - async fn is_online(&self) -> bool { - // TODO: 连接状态 - if (node_service_time_out_client(&self.addr).await).is_ok() { - return true; - } - false + fn host_name(&self) -> String { + self.endpoint.host_port() } #[tracing::instrument(skip(self))] fn endpoint(&self) -> Endpoint { @@ -100,6 +100,19 @@ impl DiskAPI for RemoteDisk { async fn close(&self) -> Result<()> { Ok(()) } + #[tracing::instrument(skip(self))] + async fn get_disk_id(&self) -> Result> { + Ok(*self.id.lock().await) + } + + #[tracing::instrument(skip(self))] + async fn set_disk_id(&self, id: Option) -> Result<()> { + let mut lock = self.id.lock().await; + *lock = id; + + Ok(()) + } + #[tracing::instrument(skip(self))] fn path(&self) -> PathBuf { self.root.clone() @@ -133,53 +146,564 @@ impl DiskAPI for RemoteDisk { } #[tracing::instrument(skip(self))] - async fn get_disk_id(&self) -> Result> { - Ok(*self.id.lock().await) - } + async fn make_volume(&self, volume: &str) -> Result<()> { + info!("make_volume"); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(MakeVolumeRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + }); - #[tracing::instrument(skip(self))] - async fn set_disk_id(&self, id: Option) -> Result<()> { - let mut lock = self.id.lock().await; - *lock = id; + let response = client.make_volume(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } Ok(()) } #[tracing::instrument(skip(self))] - async fn read_all(&self, volume: &str, path: &str) -> Result> { - info!("read_all {}/{}", volume, path); + async fn make_volumes(&self, volumes: Vec<&str>) -> Result<()> { + info!("make_volumes"); let mut client = node_service_time_out_client(&self.addr) .await .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(ReadAllRequest { + let request = Request::new(MakeVolumesRequest { disk: self.endpoint.to_string(), - volume: volume.to_string(), - path: path.to_string(), + volumes: volumes.iter().map(|s| (*s).to_string()).collect(), }); - let response = client.read_all(request).await?.into_inner(); + let response = client.make_volumes(request).await?.into_inner(); if !response.success { - return Err(Error::new(DiskError::FileNotFound)); + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; } - Ok(response.data) + Ok(()) } #[tracing::instrument(skip(self))] - async fn write_all(&self, volume: &str, path: &str, data: Vec) -> Result<()> { - info!("write_all"); + async fn list_volumes(&self) -> Result> { + info!("list_volumes"); let mut client = node_service_time_out_client(&self.addr) .await .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(WriteAllRequest { + let request = Request::new(ListVolumesRequest { + disk: self.endpoint.to_string(), + }); + + let response = client.list_volumes(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + let infos = response + .volume_infos + .into_iter() + .filter_map(|json_str| serde_json::from_str::(&json_str).ok()) + .collect(); + + Ok(infos) + } + + #[tracing::instrument(skip(self))] + async fn stat_volume(&self, volume: &str) -> Result { + info!("stat_volume"); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(StatVolumeRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + }); + + let response = client.stat_volume(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + let volume_info = serde_json::from_str::(&response.volume_info)?; + + Ok(volume_info) + } + + #[tracing::instrument(skip(self))] + async fn delete_volume(&self, volume: &str) -> Result<()> { + info!("delete_volume {}/{}", self.endpoint.to_string(), volume); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(DeleteVolumeRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + }); + + let response = client.delete_volume(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + Ok(()) + } + + // FIXME: TODO: use writer + #[tracing::instrument(skip(self, wr))] + async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()> { + let now = std::time::SystemTime::now(); + info!("walk_dir {}/{}/{:?}", self.endpoint.to_string(), opts.bucket, opts.filter_prefix); + let mut wr = wr; + let mut out = MetacacheWriter::new(&mut wr); + let mut buf = Vec::new(); + opts.serialize(&mut Serializer::new(&mut buf))?; + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(WalkDirRequest { + disk: self.endpoint.to_string(), + walk_dir_options: buf, + }); + let mut response = client.walk_dir(request).await?.into_inner(); + + loop { + match response.next().await { + Some(Ok(resp)) => { + if !resp.success { + return Err(Error::from_string(resp.error_info.unwrap_or("".to_string()))); + } + let entry = serde_json::from_str::(&resp.meta_cache_entry) + .map_err(|_| Error::from_string(format!("Unexpected response: {:?}", response)))?; + out.write_obj(&entry).await?; + } + None => break, + _ => return Err(Error::from_string(format!("Unexpected response: {:?}", response))), + } + } + + info!( + "walk_dir {}/{:?} done {:?}", + opts.bucket, + opts.filter_prefix, + now.elapsed().unwrap_or_default() + ); + Ok(()) + } + + #[tracing::instrument(skip(self))] + async fn delete_version( + &self, + volume: &str, + path: &str, + fi: FileInfo, + force_del_marker: bool, + opts: DeleteOptions, + ) -> Result<()> { + info!("delete_version"); + let file_info = serde_json::to_string(&fi)?; + let opts = serde_json::to_string(&opts)?; + + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(DeleteVersionRequest { disk: self.endpoint.to_string(), volume: volume.to_string(), path: path.to_string(), - data, + file_info, + force_del_marker, + opts, }); - let response = client.write_all(request).await?.into_inner(); + let response = client.delete_version(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + // let raw_file_info = serde_json::from_str::(&response.raw_file_info)?; + + Ok(()) + } + + #[tracing::instrument(skip(self))] + async fn delete_versions( + &self, + volume: &str, + versions: Vec, + opts: DeleteOptions, + ) -> Result>> { + info!("delete_versions"); + let opts = serde_json::to_string(&opts)?; + let mut versions_str = Vec::with_capacity(versions.len()); + for file_info_versions in versions.iter() { + versions_str.push(serde_json::to_string(file_info_versions)?); + } + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(DeleteVersionsRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + versions: versions_str, + opts, + }); + + let response = client.delete_versions(request).await?.into_inner(); + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + let errors = response + .errors + .iter() + .map(|error| { + if error.is_empty() { + None + } else { + Some(Error::from_string(error)) + } + }) + .collect(); + + Ok(errors) + } + + #[tracing::instrument(skip(self))] + async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()> { + info!("delete_paths"); + let paths = paths.to_owned(); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(DeletePathsRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + paths, + }); + + let response = client.delete_paths(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + Ok(()) + } + + #[tracing::instrument(skip(self))] + async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { + info!("write_metadata {}/{}", volume, path); + let file_info = serde_json::to_string(&fi)?; + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(WriteMetadataRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + file_info, + }); + + let response = client.write_metadata(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + Ok(()) + } + + #[tracing::instrument(skip(self))] + async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> { + info!("update_metadata"); + let file_info = serde_json::to_string(&fi)?; + let opts = serde_json::to_string(&opts)?; + + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(UpdateMetadataRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + file_info, + opts, + }); + + let response = client.update_metadata(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + Ok(()) + } + + #[tracing::instrument(skip(self))] + async fn read_version( + &self, + _org_volume: &str, + volume: &str, + path: &str, + version_id: &str, + opts: &ReadOptions, + ) -> Result { + info!("read_version"); + let opts = serde_json::to_string(opts)?; + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(ReadVersionRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + version_id: version_id.to_string(), + opts, + }); + + let response = client.read_version(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + let file_info = serde_json::from_str::(&response.file_info)?; + + Ok(file_info) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result { + info!("read_xl {}/{}/{}", self.endpoint.to_string(), volume, path); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(ReadXlRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + read_data, + }); + + let response = client.read_xl(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + let raw_file_info = serde_json::from_str::(&response.raw_file_info)?; + + Ok(raw_file_info) + } + + #[tracing::instrument(skip(self))] + async fn rename_data( + &self, + src_volume: &str, + src_path: &str, + fi: FileInfo, + dst_volume: &str, + dst_path: &str, + ) -> Result { + info!("rename_data {}/{}/{}/{}", self.addr, self.endpoint.to_string(), dst_volume, dst_path); + let file_info = serde_json::to_string(&fi)?; + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(RenameDataRequest { + disk: self.endpoint.to_string(), + src_volume: src_volume.to_string(), + src_path: src_path.to_string(), + file_info, + dst_volume: dst_volume.to_string(), + dst_path: dst_path.to_string(), + }); + + let response = client.rename_data(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + let rename_data_resp = serde_json::from_str::(&response.rename_data_resp)?; + + Ok(rename_data_resp) + } + + #[tracing::instrument(skip(self))] + async fn list_dir(&self, _origvolume: &str, volume: &str, _dir_path: &str, _count: i32) -> Result> { + info!("list_dir {}/{}", volume, _dir_path); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(ListDirRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + }); + + let response = client.list_dir(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + Ok(response.volumes) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn read_file(&self, volume: &str, path: &str) -> Result { + info!("read_file {}/{}", volume, path); + Ok(Box::new( + HttpFileReader::new(self.endpoint.grid_host().as_str(), self.endpoint.to_string().as_str(), volume, path, 0, 0) + .await?, + )) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn read_file_stream(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result { + info!("read_file_stream {}/{}/{}", self.endpoint.to_string(), volume, path); + Ok(Box::new( + HttpFileReader::new( + self.endpoint.grid_host().as_str(), + self.endpoint.to_string().as_str(), + volume, + path, + offset, + length, + ) + .await?, + )) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn append_file(&self, volume: &str, path: &str) -> Result { + info!("append_file {}/{}", volume, path); + Ok(Box::new(HttpFileWriter::new( + self.endpoint.grid_host().as_str(), + self.endpoint.to_string().as_str(), + volume, + path, + 0, + true, + )?)) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, file_size: usize) -> Result { + info!("create_file {}/{}/{}", self.endpoint.to_string(), volume, path); + Ok(Box::new(HttpFileWriter::new( + self.endpoint.grid_host().as_str(), + self.endpoint.to_string().as_str(), + volume, + path, + file_size, + false, + )?)) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> { + info!("rename_file"); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(RenameFileRequst { + disk: self.endpoint.to_string(), + src_volume: src_volume.to_string(), + src_path: src_path.to_string(), + dst_volume: dst_volume.to_string(), + dst_path: dst_path.to_string(), + }); + + let response = client.rename_file(request).await?.into_inner(); + + if !response.success { + return if let Some(err) = &response.error { + Err(proto_err_to_err(err)) + } else { + Err(Error::from_string("")) + }; + } + + Ok(()) + } + + #[tracing::instrument(skip(self))] + async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec) -> Result<()> { + info!("rename_part {}/{}", src_volume, src_path); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(RenamePartRequst { + disk: self.endpoint.to_string(), + src_volume: src_volume.to_string(), + src_path: src_path.to_string(), + dst_volume: dst_volume.to_string(), + dst_path: dst_path.to_string(), + meta, + }); + + let response = client.rename_part(request).await?.into_inner(); if !response.success { return if let Some(err) = &response.error { @@ -277,553 +801,6 @@ impl DiskAPI for RemoteDisk { Ok(check_parts_resp) } - #[tracing::instrument(skip(self))] - async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec) -> Result<()> { - info!("rename_part {}/{}", src_volume, src_path); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(RenamePartRequst { - disk: self.endpoint.to_string(), - src_volume: src_volume.to_string(), - src_path: src_path.to_string(), - dst_volume: dst_volume.to_string(), - dst_path: dst_path.to_string(), - meta, - }); - - let response = client.rename_part(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> { - info!("rename_file"); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(RenameFileRequst { - disk: self.endpoint.to_string(), - src_volume: src_volume.to_string(), - src_path: src_path.to_string(), - dst_volume: dst_volume.to_string(), - dst_path: dst_path.to_string(), - }); - - let response = client.rename_file(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, file_size: usize) -> Result { - info!("create_file {}/{}/{}", self.endpoint.to_string(), volume, path); - Ok(Box::new(HttpFileWriter::new( - self.endpoint.grid_host().as_str(), - self.endpoint.to_string().as_str(), - volume, - path, - file_size, - false, - )?)) - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn append_file(&self, volume: &str, path: &str) -> Result { - info!("append_file {}/{}", volume, path); - Ok(Box::new(HttpFileWriter::new( - self.endpoint.grid_host().as_str(), - self.endpoint.to_string().as_str(), - volume, - path, - 0, - true, - )?)) - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn read_file(&self, volume: &str, path: &str) -> Result { - info!("read_file {}/{}", volume, path); - Ok(Box::new( - HttpFileReader::new(self.endpoint.grid_host().as_str(), self.endpoint.to_string().as_str(), volume, path, 0, 0) - .await?, - )) - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn read_file_stream(&self, volume: &str, path: &str, offset: usize, length: usize) -> Result { - info!("read_file_stream {}/{}/{}", self.endpoint.to_string(), volume, path); - Ok(Box::new( - HttpFileReader::new( - self.endpoint.grid_host().as_str(), - self.endpoint.to_string().as_str(), - volume, - path, - offset, - length, - ) - .await?, - )) - } - - #[tracing::instrument(skip(self))] - async fn list_dir(&self, _origvolume: &str, volume: &str, _dir_path: &str, _count: i32) -> Result> { - info!("list_dir {}/{}", volume, _dir_path); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(ListDirRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - }); - - let response = client.list_dir(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(response.volumes) - } - - // FIXME: TODO: use writer - #[tracing::instrument(skip(self, wr))] - async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()> { - let now = std::time::SystemTime::now(); - info!("walk_dir {}/{}/{:?}", self.endpoint.to_string(), opts.bucket, opts.filter_prefix); - let mut wr = wr; - let mut out = MetacacheWriter::new(&mut wr); - let mut buf = Vec::new(); - opts.serialize(&mut Serializer::new(&mut buf))?; - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(WalkDirRequest { - disk: self.endpoint.to_string(), - walk_dir_options: buf, - }); - let mut response = client.walk_dir(request).await?.into_inner(); - - loop { - match response.next().await { - Some(Ok(resp)) => { - if !resp.success { - return Err(Error::from_string(resp.error_info.unwrap_or("".to_string()))); - } - let entry = serde_json::from_str::(&resp.meta_cache_entry) - .map_err(|_| Error::from_string(format!("Unexpected response: {:?}", response)))?; - out.write_obj(&entry).await?; - } - None => break, - _ => return Err(Error::from_string(format!("Unexpected response: {:?}", response))), - } - } - - info!( - "walk_dir {}/{:?} done {:?}", - opts.bucket, - opts.filter_prefix, - now.elapsed().unwrap_or_default() - ); - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn rename_data( - &self, - src_volume: &str, - src_path: &str, - fi: FileInfo, - dst_volume: &str, - dst_path: &str, - ) -> Result { - info!("rename_data {}/{}/{}/{}", self.addr, self.endpoint.to_string(), dst_volume, dst_path); - let file_info = serde_json::to_string(&fi)?; - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(RenameDataRequest { - disk: self.endpoint.to_string(), - src_volume: src_volume.to_string(), - src_path: src_path.to_string(), - file_info, - dst_volume: dst_volume.to_string(), - dst_path: dst_path.to_string(), - }); - - let response = client.rename_data(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - let rename_data_resp = serde_json::from_str::(&response.rename_data_resp)?; - - Ok(rename_data_resp) - } - - #[tracing::instrument(skip(self))] - async fn make_volumes(&self, volumes: Vec<&str>) -> Result<()> { - info!("make_volumes"); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(MakeVolumesRequest { - disk: self.endpoint.to_string(), - volumes: volumes.iter().map(|s| (*s).to_string()).collect(), - }); - - let response = client.make_volumes(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn make_volume(&self, volume: &str) -> Result<()> { - info!("make_volume"); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(MakeVolumeRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - }); - - let response = client.make_volume(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn list_volumes(&self) -> Result> { - info!("list_volumes"); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(ListVolumesRequest { - disk: self.endpoint.to_string(), - }); - - let response = client.list_volumes(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - let infos = response - .volume_infos - .into_iter() - .filter_map(|json_str| serde_json::from_str::(&json_str).ok()) - .collect(); - - Ok(infos) - } - - #[tracing::instrument(skip(self))] - async fn stat_volume(&self, volume: &str) -> Result { - info!("stat_volume"); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(StatVolumeRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - }); - - let response = client.stat_volume(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - let volume_info = serde_json::from_str::(&response.volume_info)?; - - Ok(volume_info) - } - - #[tracing::instrument(skip(self))] - async fn delete_paths(&self, volume: &str, paths: &[String]) -> Result<()> { - info!("delete_paths"); - let paths = paths.to_owned(); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(DeletePathsRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - paths, - }); - - let response = client.delete_paths(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: &UpdateMetadataOpts) -> Result<()> { - info!("update_metadata"); - let file_info = serde_json::to_string(&fi)?; - let opts = serde_json::to_string(&opts)?; - - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(UpdateMetadataRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - path: path.to_string(), - file_info, - opts, - }); - - let response = client.update_metadata(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { - info!("write_metadata {}/{}", volume, path); - let file_info = serde_json::to_string(&fi)?; - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(WriteMetadataRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - path: path.to_string(), - file_info, - }); - - let response = client.write_metadata(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn read_version( - &self, - _org_volume: &str, - volume: &str, - path: &str, - version_id: &str, - opts: &ReadOptions, - ) -> Result { - info!("read_version"); - let opts = serde_json::to_string(opts)?; - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(ReadVersionRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - path: path.to_string(), - version_id: version_id.to_string(), - opts, - }); - - let response = client.read_version(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - let file_info = serde_json::from_str::(&response.file_info)?; - - Ok(file_info) - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result { - info!("read_xl {}/{}/{}", self.endpoint.to_string(), volume, path); - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(ReadXlRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - path: path.to_string(), - read_data, - }); - - let response = client.read_xl(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - let raw_file_info = serde_json::from_str::(&response.raw_file_info)?; - - Ok(raw_file_info) - } - - #[tracing::instrument(skip(self))] - async fn delete_version( - &self, - volume: &str, - path: &str, - fi: FileInfo, - force_del_marker: bool, - opts: DeleteOptions, - ) -> Result<()> { - info!("delete_version"); - let file_info = serde_json::to_string(&fi)?; - let opts = serde_json::to_string(&opts)?; - - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(DeleteVersionRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - path: path.to_string(), - file_info, - force_del_marker, - opts, - }); - - let response = client.delete_version(request).await?.into_inner(); - - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - - // let raw_file_info = serde_json::from_str::(&response.raw_file_info)?; - - Ok(()) - } - - #[tracing::instrument(skip(self))] - async fn delete_versions( - &self, - volume: &str, - versions: Vec, - opts: DeleteOptions, - ) -> Result>> { - info!("delete_versions"); - let opts = serde_json::to_string(&opts)?; - let mut versions_str = Vec::with_capacity(versions.len()); - for file_info_versions in versions.iter() { - versions_str.push(serde_json::to_string(file_info_versions)?); - } - let mut client = node_service_time_out_client(&self.addr) - .await - .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(DeleteVersionsRequest { - disk: self.endpoint.to_string(), - volume: volume.to_string(), - versions: versions_str, - opts, - }); - - let response = client.delete_versions(request).await?.into_inner(); - if !response.success { - return if let Some(err) = &response.error { - Err(proto_err_to_err(err)) - } else { - Err(Error::from_string("")) - }; - } - let errors = response - .errors - .iter() - .map(|error| { - if error.is_empty() { - None - } else { - Some(Error::from_string(error)) - } - }) - .collect(); - - Ok(errors) - } - #[tracing::instrument(skip(self))] async fn read_multiple(&self, req: ReadMultipleReq) -> Result> { info!("read_multiple {}/{}/{}", self.endpoint.to_string(), req.bucket, req.prefix); @@ -856,17 +833,19 @@ impl DiskAPI for RemoteDisk { } #[tracing::instrument(skip(self))] - async fn delete_volume(&self, volume: &str) -> Result<()> { - info!("delete_volume {}/{}", self.endpoint.to_string(), volume); + async fn write_all(&self, volume: &str, path: &str, data: Vec) -> Result<()> { + info!("write_all"); let mut client = node_service_time_out_client(&self.addr) .await .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; - let request = Request::new(DeleteVolumeRequest { + let request = Request::new(WriteAllRequest { disk: self.endpoint.to_string(), volume: volume.to_string(), + path: path.to_string(), + data, }); - let response = client.delete_volume(request).await?.into_inner(); + let response = client.write_all(request).await?.into_inner(); if !response.success { return if let Some(err) = &response.error { @@ -879,6 +858,27 @@ impl DiskAPI for RemoteDisk { Ok(()) } + #[tracing::instrument(skip(self))] + async fn read_all(&self, volume: &str, path: &str) -> Result> { + info!("read_all {}/{}", volume, path); + let mut client = node_service_time_out_client(&self.addr) + .await + .map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; + let request = Request::new(ReadAllRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + }); + + let response = client.read_all(request).await?.into_inner(); + + if !response.success { + return Err(Error::new(DiskError::FileNotFound)); + } + + Ok(response.data) + } + #[tracing::instrument(skip(self))] async fn disk_info(&self, opts: &DiskInfoOptions) -> Result { let opts = serde_json::to_string(&opts)?; diff --git a/ecstore/src/disks_layout.rs b/ecstore/src/disks_layout.rs index 75eab78ec..703ae18f8 100644 --- a/ecstore/src/disks_layout.rs +++ b/ecstore/src/disks_layout.rs @@ -94,13 +94,10 @@ impl DisksLayout { let is_ellipses = args.iter().any(|v| has_ellipses(&[v])); - let set_drive_count_env = match env::var(ENV_RUSTFS_ERASURE_SET_DRIVE_COUNT) { - Ok(res) => res, - Err(err) => { - debug!("{} not set use default:0, {:?}", ENV_RUSTFS_ERASURE_SET_DRIVE_COUNT, err); - "0".to_string() - } - }; + let set_drive_count_env = env::var(ENV_RUSTFS_ERASURE_SET_DRIVE_COUNT).unwrap_or_else(|err| { + debug!("{} not set use default:0, {:?}", ENV_RUSTFS_ERASURE_SET_DRIVE_COUNT, err); + "0".to_string() + }); let set_drive_count: usize = set_drive_count_env.parse()?; // None of the args have ellipses use the old style. @@ -290,7 +287,7 @@ impl EndpointSet { } } -/// returns a greatest common divisor of all the ellipses sizes. +/// returns the greatest common divisor of all the ellipses sizes. fn get_divisible_size(total_sizes: &[usize]) -> usize { fn gcd(mut x: usize, mut y: usize) -> usize { while y != 0 { @@ -446,7 +443,6 @@ fn get_total_sizes(arg_patterns: &[ArgPattern]) -> Vec { #[cfg(test)] mod test { - use super::*; impl PartialEq for EndpointSet { @@ -868,7 +864,7 @@ mod test { }, success: true, }, - // More than 1 ellipses per argument for standalone setup. + // More than an ellipse per argument for standalone setup. TestCase { num: 15, arg: "/export{1...10}/disk{1...10}", diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index 7329e77cd..0edc97c5c 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -113,7 +113,7 @@ impl Erasure { let blocks_inner = blocks.clone(); async move { if let Some(w) = w_op { - (w.write(blocks_inner[i_inner].clone()).await).err() + w.write(blocks_inner[i_inner].clone()).await.err() } else { Some(Error::new(DiskError::DiskNotFound)) } diff --git a/ecstore/src/heal/data_scanner.rs b/ecstore/src/heal/data_scanner.rs index c806267f2..627f044e6 100644 --- a/ecstore/src/heal/data_scanner.rs +++ b/ecstore/src/heal/data_scanner.rs @@ -151,7 +151,7 @@ pub async fn init_data_scanner() { let mut r = rand::thread_rng(); r.gen_range(0.0..1.0) }; - let duration = Duration::from_secs_f64(random * (SCANNER_CYCLE.load(std::sync::atomic::Ordering::SeqCst) as f64)); + let duration = Duration::from_secs_f64(random * (SCANNER_CYCLE.load(Ordering::SeqCst) as f64)); let sleep_duration = if duration < Duration::new(1, 0) { Duration::new(1, 0) } else { @@ -227,7 +227,7 @@ async fn run_data_scanner() { } } stop_fn(&res).await; - sleep(Duration::from_secs(SCANNER_CYCLE.load(std::sync::atomic::Ordering::SeqCst))).await; + sleep(Duration::from_secs(SCANNER_CYCLE.load(Ordering::SeqCst))).await; } } @@ -383,7 +383,7 @@ impl CurrentScannerCycle { Deserialize::deserialize(&mut Deserializer::new(&buf[..])).expect("Deserialization failed"); self.cycle_completed = u; } - name => return Err(Error::msg(format!("not suport field name {}", name))), + name => return Err(Error::msg(format!("not support field name {}", name))), } } @@ -435,8 +435,8 @@ impl ScannerItem { path_join(&[PathBuf::from(self.prefix.clone()), PathBuf::from(self.object_name.clone())]) } - pub async fn apply_versions_actions(&self, fivs: &[FileInfo]) -> Result> { - let obj_infos = self.apply_newer_noncurrent_version_limit(fivs).await?; + pub async fn apply_versions_actions(&self, fives: &[FileInfo]) -> Result> { + let obj_infos = self.apply_newer_noncurrent_version_limit(fives).await?; if obj_infos.len() >= SCANNER_EXCESS_OBJECT_VERSIONS.load(Ordering::SeqCst) as usize { // todo } @@ -453,16 +453,16 @@ impl ScannerItem { Ok(obj_infos) } - pub async fn apply_newer_noncurrent_version_limit(&self, fivs: &[FileInfo]) -> Result> { + pub async fn apply_newer_noncurrent_version_limit(&self, fives: &[FileInfo]) -> Result> { // let done = ScannerMetrics::time(ScannerMetric::ApplyNonCurrent); let versioned = match BucketVersioningSys::get(&self.bucket).await { Ok(vcfg) => vcfg.versioned(self.object_path().to_str().unwrap_or_default()), Err(_) => false, }; - let mut object_infos = Vec::with_capacity(fivs.len()); + let mut object_infos = Vec::with_capacity(fives.len()); if self.lifecycle.is_none() { - for info in fivs.iter() { + for info in fives.iter() { object_infos.push(info.to_object_info(&self.bucket, &self.object_path().to_string_lossy(), versioned)); } return Ok(object_infos); @@ -1000,7 +1000,7 @@ impl FolderScanner { if !into.compacted { self.new_cache.reduce_children_of( &this_hash, - DATA_SCANNER_COMPACT_AT_CHILDREN.try_into().unwrap(), + DATA_SCANNER_COMPACT_AT_CHILDREN.try_into()?, self.new_cache.info.name != folder.name, ); } diff --git a/ecstore/src/heal/data_usage_cache.rs b/ecstore/src/heal/data_usage_cache.rs index 50a570f5e..b22eda382 100644 --- a/ecstore/src/heal/data_usage_cache.rs +++ b/ecstore/src/heal/data_usage_cache.rs @@ -241,7 +241,7 @@ impl ReplicationAllStats { #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct DataUsageEntry { pub children: DataUsageHashMap, - // These fields do no include any children. + // These fields do not include any children. pub size: usize, pub objects: usize, pub versions: usize, diff --git a/ecstore/src/io.rs b/ecstore/src/io.rs index 2bd02c117..185918b43 100644 --- a/ecstore/src/io.rs +++ b/ecstore/src/io.rs @@ -27,14 +27,14 @@ pub const READ_BUFFER_SIZE: usize = 1024 * 1024; #[derive(Debug)] pub struct HttpFileWriter { wd: tokio::io::DuplexStream, - err_rx: oneshot::Receiver, + err_rx: oneshot::Receiver, } impl HttpFileWriter { - pub fn new(url: &str, disk: &str, volume: &str, path: &str, size: usize, append: bool) -> std::io::Result { + pub fn new(url: &str, disk: &str, volume: &str, path: &str, size: usize, append: bool) -> io::Result { let (rd, wd) = tokio::io::duplex(READ_BUFFER_SIZE); - let (err_tx, err_rx) = oneshot::channel::(); + let (err_tx, err_rx) = oneshot::channel::(); let body = reqwest::Body::wrap_stream(ReaderStream::with_capacity(rd, READ_BUFFER_SIZE)); @@ -58,7 +58,7 @@ impl HttpFileWriter { .body(body) .send() .await - .map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e)) + .map_err(|e| io::Error::new(io::ErrorKind::Other, e)) { error!("HttpFileWriter put file err: {:?}", err); @@ -74,11 +74,7 @@ impl HttpFileWriter { impl AsyncWrite for HttpFileWriter { #[tracing::instrument(level = "debug", skip(self, buf))] - fn poll_write( - mut self: Pin<&mut Self>, - cx: &mut std::task::Context<'_>, - buf: &[u8], - ) -> Poll> { + fn poll_write(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll> { if let Ok(err) = self.as_mut().err_rx.try_recv() { return Poll::Ready(Err(err)); } @@ -87,12 +83,12 @@ impl AsyncWrite for HttpFileWriter { } #[tracing::instrument(level = "debug", skip(self))] - fn poll_flush(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll> { + fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { Pin::new(&mut self.wd).poll_flush(cx) } #[tracing::instrument(level = "debug", skip(self))] - fn poll_shutdown(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll> { + fn poll_shutdown(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { Pin::new(&mut self.wd).poll_shutdown(cx) } } @@ -102,7 +98,7 @@ pub struct HttpFileReader { } impl HttpFileReader { - pub async fn new(url: &str, disk: &str, volume: &str, path: &str, offset: usize, length: usize) -> std::io::Result { + pub async fn new(url: &str, disk: &str, volume: &str, path: &str, offset: usize, length: usize) -> io::Result { let resp = reqwest::Client::new() .get(format!( "{}/rustfs/rpc/read_file_stream?disk={}&volume={}&path={}&offset={}&length={}", @@ -115,16 +111,16 @@ impl HttpFileReader { )) .send() .await - .map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?; + .map_err(|e| io::Error::new(io::ErrorKind::Other, e))?; - let inner = Box::new(StreamReader::new(resp.bytes_stream().map_err(std::io::Error::other))); + let inner = Box::new(StreamReader::new(resp.bytes_stream().map_err(io::Error::other))); Ok(Self { inner }) } } impl AsyncRead for HttpFileReader { - fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { Pin::new(&mut self.inner).poll_read(cx, buf) } } @@ -172,7 +168,7 @@ impl Etag for EtagReader { impl AsyncRead for EtagReader { #[tracing::instrument(level = "info", skip_all)] - fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { let me = self.project(); loop { diff --git a/ecstore/src/notification_sys.rs b/ecstore/src/notification_sys.rs index ed4649069..cf538e24d 100644 --- a/ecstore/src/notification_sys.rs +++ b/ecstore/src/notification_sys.rs @@ -187,7 +187,7 @@ fn get_offline_disks(offline_host: &str, endpoints: &EndpointServerPools) -> Vec if (offline_host.is_empty() && ep.is_local) || offline_host == ep.host_port() { offline_disks.push(madmin::Disk { endpoint: ep.to_string(), - state: madmin::ItemState::Offline.to_string().to_owned(), + state: ItemState::Offline.to_string().to_owned(), pool_index: ep.pool_idx, set_index: ep.set_idx, disk_index: ep.disk_idx, diff --git a/ecstore/src/peer_rest_client.rs b/ecstore/src/peer_rest_client.rs index 9a121c9b3..5782990e5 100644 --- a/ecstore/src/peer_rest_client.rs +++ b/ecstore/src/peer_rest_client.rs @@ -89,7 +89,7 @@ impl PeerRestClient { let data = response.storage_info; let mut buf = Deserializer::new(Cursor::new(data)); - let storage_info: madmin::StorageInfo = Deserialize::deserialize(&mut buf).unwrap(); + let storage_info: madmin::StorageInfo = Deserialize::deserialize(&mut buf)?; Ok(storage_info) } @@ -110,7 +110,7 @@ impl PeerRestClient { let data = response.server_properties; let mut buf = Deserializer::new(Cursor::new(data)); - let storage_properties: ServerProperties = Deserialize::deserialize(&mut buf).unwrap(); + let storage_properties: ServerProperties = Deserialize::deserialize(&mut buf)?; Ok(storage_properties) } @@ -131,7 +131,7 @@ impl PeerRestClient { let data = response.cpus; let mut buf = Deserializer::new(Cursor::new(data)); - let cpus: Cpus = Deserialize::deserialize(&mut buf).unwrap(); + let cpus: Cpus = Deserialize::deserialize(&mut buf)?; Ok(cpus) } @@ -152,7 +152,7 @@ impl PeerRestClient { let data = response.net_info; let mut buf = Deserializer::new(Cursor::new(data)); - let net_info: NetInfo = Deserialize::deserialize(&mut buf).unwrap(); + let net_info: NetInfo = Deserialize::deserialize(&mut buf)?; Ok(net_info) } @@ -173,7 +173,7 @@ impl PeerRestClient { let data = response.partitions; let mut buf = Deserializer::new(Cursor::new(data)); - let partitions: Partitions = Deserialize::deserialize(&mut buf).unwrap(); + let partitions: Partitions = Deserialize::deserialize(&mut buf)?; Ok(partitions) } @@ -194,7 +194,7 @@ impl PeerRestClient { let data = response.os_info; let mut buf = Deserializer::new(Cursor::new(data)); - let os_info: OsInfo = Deserialize::deserialize(&mut buf).unwrap(); + let os_info: OsInfo = Deserialize::deserialize(&mut buf)?; Ok(os_info) } @@ -215,7 +215,7 @@ impl PeerRestClient { let data = response.sys_services; let mut buf = Deserializer::new(Cursor::new(data)); - let sys_services: SysService = Deserialize::deserialize(&mut buf).unwrap(); + let sys_services: SysService = Deserialize::deserialize(&mut buf)?; Ok(sys_services) } @@ -236,7 +236,7 @@ impl PeerRestClient { let data = response.sys_config; let mut buf = Deserializer::new(Cursor::new(data)); - let sys_config: SysConfig = Deserialize::deserialize(&mut buf).unwrap(); + let sys_config: SysConfig = Deserialize::deserialize(&mut buf)?; Ok(sys_config) } @@ -257,7 +257,7 @@ impl PeerRestClient { let data = response.sys_errors; let mut buf = Deserializer::new(Cursor::new(data)); - let sys_errors: SysErrors = Deserialize::deserialize(&mut buf).unwrap(); + let sys_errors: SysErrors = Deserialize::deserialize(&mut buf)?; Ok(sys_errors) } @@ -278,7 +278,7 @@ impl PeerRestClient { let data = response.mem_info; let mut buf = Deserializer::new(Cursor::new(data)); - let mem_info: MemInfo = Deserialize::deserialize(&mut buf).unwrap(); + let mem_info: MemInfo = Deserialize::deserialize(&mut buf)?; Ok(mem_info) } @@ -306,7 +306,7 @@ impl PeerRestClient { let data = response.realtime_metrics; let mut buf = Deserializer::new(Cursor::new(data)); - let realtime_metrics: RealtimeMetrics = Deserialize::deserialize(&mut buf).unwrap(); + let realtime_metrics: RealtimeMetrics = Deserialize::deserialize(&mut buf)?; Ok(realtime_metrics) } @@ -327,7 +327,7 @@ impl PeerRestClient { let data = response.proc_info; let mut buf = Deserializer::new(Cursor::new(data)); - let proc_info: ProcInfo = Deserialize::deserialize(&mut buf).unwrap(); + let proc_info: ProcInfo = Deserialize::deserialize(&mut buf)?; Ok(proc_info) } @@ -603,20 +603,20 @@ impl PeerRestClient { let data = response.bg_heal_state; let mut buf = Deserializer::new(Cursor::new(data)); - let bg_heal_state: BgHealState = Deserialize::deserialize(&mut buf).unwrap(); + let bg_heal_state: BgHealState = Deserialize::deserialize(&mut buf)?; Ok(bg_heal_state) } pub async fn get_metacache_listing(&self) -> Result<()> { - let mut _client = node_service_time_out_client(&self.grid_host) + let _client = node_service_time_out_client(&self.grid_host) .await .map_err(|err| Error::msg(err.to_string()))?; todo!() } pub async fn update_metacache_listing(&self) -> Result<()> { - let mut _client = node_service_time_out_client(&self.grid_host) + let _client = node_service_time_out_client(&self.grid_host) .await .map_err(|err| Error::msg(err.to_string()))?; todo!() diff --git a/ecstore/src/pools.rs b/ecstore/src/pools.rs index 97615a6a2..f59ee79b9 100644 --- a/ecstore/src/pools.rs +++ b/ecstore/src/pools.rs @@ -127,7 +127,7 @@ impl PoolMeta { } let mut buf = Deserializer::new(Cursor::new(&data[4..])); - let meta: PoolMeta = Deserialize::deserialize(&mut buf).unwrap(); + let meta: PoolMeta = Deserialize::deserialize(&mut buf)?; *self = meta; if self.version != POOL_META_VERSION { @@ -141,8 +141,8 @@ impl PoolMeta { return Ok(()); } let mut data = Vec::new(); - data.write_u16::(POOL_META_FORMAT).unwrap(); - data.write_u16::(POOL_META_VERSION).unwrap(); + data.write_u16::(POOL_META_FORMAT)?; + data.write_u16::(POOL_META_VERSION)?; let mut buf = Vec::new(); self.serialize(&mut Serializer::new(&mut buf))?; data.write_all(&buf)?; diff --git a/ecstore/src/rebalance.rs b/ecstore/src/rebalance.rs index 1f95042fd..349c56040 100644 --- a/ecstore/src/rebalance.rs +++ b/ecstore/src/rebalance.rs @@ -139,7 +139,7 @@ pub struct DiskStat { #[derive(Debug, Default, Serialize, Deserialize)] pub struct RebalanceMeta { #[serde(skip)] - pub cancel: Option>, // To be invoked on rebalance-stop + pub cancel: Option>, // To be invoked on rebalance-stop #[serde(skip)] pub last_refreshed_at: Option, #[serde(rename = "stopTs")] diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 9ffe1f5ea..83b00b403 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -36,7 +36,6 @@ use crate::{ store_err::{is_err_object_not_found, to_object_err, StorageError}, store_init::{load_format_erasure, ErasureError}, utils::{ - self, crypto::{base64_decode, base64_encode, hex}, path::{encode_dir_object, has_suffix, SLASH_SEPARATOR}, }, @@ -174,7 +173,7 @@ impl SetDisks { let mut rng = thread_rng(); disks.shuffle(&mut rng); - numbers.shuffle(&mut rand::thread_rng()); + numbers.shuffle(&mut thread_rng()); } for &i in numbers.iter() { @@ -395,7 +394,7 @@ impl SetDisks { } fn reduce_common_data_dir(data_dirs: &Vec>, write_quorum: usize) -> Option { - let mut data_dirs_count = std::collections::HashMap::new(); + let mut data_dirs_count = HashMap::new(); for ddir in data_dirs { *data_dirs_count.entry(ddir).or_insert(0) += 1; @@ -454,10 +453,7 @@ impl SetDisks { let errs: Vec> = join_all(futures) .await .into_iter() - .map(|e| match e { - Ok(e) => e, - Err(_) => Some(Error::new(DiskError::Unexpected)), - }) + .map(|e| e.unwrap_or_else(|_| Some(Error::new(DiskError::Unexpected)))) .collect(); if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) { @@ -726,7 +722,7 @@ impl SetDisks { } if max_occ == 0 { - // Did not found anything useful + // Did not find anything useful return -1; } cparity @@ -1039,7 +1035,7 @@ impl SetDisks { hasher.flush()?; - meta_hashs[i] = Some(utils::crypto::hex(hasher.clone().finalize().as_slice())); + meta_hashs[i] = Some(hex(hasher.clone().finalize().as_slice())); hasher.reset(); } @@ -2331,7 +2327,7 @@ impl SetDisks { // Allow for dangling deletes, on versions that have DataDir missing etc. // this would end up restoring the correct readable versions. - match self + return match self .delete_if_dang_ling( bucket, object, @@ -2355,10 +2351,7 @@ impl SetDisks { for _ in 0..errs.len() { t_errs.push(None); } - return Ok(( - self.default_heal_result(m, &t_errs, bucket, object, version_id).await, - Some(derr), - )); + Ok((self.default_heal_result(m, &t_errs, bucket, object, version_id).await, Some(derr))) } Err(err) => { // t_errs = vec![Some(err.clone()); errs.len()]; @@ -2367,13 +2360,13 @@ impl SetDisks { t_errs.push(Some(clone_err(&err))); } - return Ok(( + Ok(( self.default_heal_result(FileInfo::default(), &t_errs, bucket, object, version_id) .await, Some(err), - )); + )) } - } + }; } if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != available_disks.len() { @@ -2908,6 +2901,7 @@ impl SetDisks { ) -> Result<()> { info!("ns_scanner"); if buckets.is_empty() { + info!("data-scanner: no buckets to scan, skipping scanner cycle"); return Ok(()); } @@ -2930,7 +2924,7 @@ impl SetDisks { // Put all buckets into channel. let (bucket_tx, bucket_rx) = mpsc::channel(buckets.len()); // Shuffle buckets to ensure total randomness of buckets, being scanned. - // Otherwise same set of buckets get scanned across erasure sets always. + // Otherwise, same set of buckets get scanned across erasure sets always. // at any given point in time. This allows different buckets to be scanned // in different order per erasure set, this wider spread is needed when // there are lots of buckets with different order of objects in them. @@ -5120,7 +5114,7 @@ impl StorageAPI for SetDisks { _ => {} } } - return Ok((result, err)); + Ok((result, err)) } #[tracing::instrument(skip(self))] @@ -5420,7 +5414,7 @@ async fn disks_with_all_parts( if let Some(data) = &meta.data { let checksum_info = meta.erasure.get_checksum_info(meta.parts[0].number); let data_len = data.len(); - let verify_err = (bitrot_verify( + let verify_err = bitrot_verify( Box::new(Cursor::new(data.clone())), data_len, meta.erasure.shard_file_size(meta.size), @@ -5428,8 +5422,8 @@ async fn disks_with_all_parts( checksum_info.hash, meta.erasure.shard_size(meta.erasure.block_size), ) - .await) - .err(); + .await + .err(); if let Some(vec) = data_errs_by_part.get_mut(&0) { if index < vec.len() { diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index 22ed4471d..5c978168a 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -121,18 +121,15 @@ impl Sets { disk = local_disk; } - let has_disk_id = match disk.as_ref().unwrap().get_disk_id().await { - Ok(res) => res, - Err(err) => { - if is_unformatted_disk(&err) { - error!("get_disk_id err {:?}", err); - } else { - warn!("get_disk_id err {:?}", err); - } - - None + let has_disk_id = disk.as_ref().unwrap().get_disk_id().await.unwrap_or_else(|err| { + if is_unformatted_disk(&err) { + error!("get_disk_id err {:?}", err); + } else { + warn!("get_disk_id err {:?}", err); } - }; + + None + }); if let Some(_disk_id) = has_disk_id { set_drive.push(disk); @@ -354,19 +351,56 @@ impl StorageAPI for Sets { } } #[tracing::instrument(skip(self))] - async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { + async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> { unimplemented!() } #[tracing::instrument(skip(self))] - async fn make_bucket(&self, _bucket: &str, _opts: &MakeBucketOptions) -> Result<()> { + async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { unimplemented!() } #[tracing::instrument(skip(self))] - async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { + async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { unimplemented!() } + #[tracing::instrument(skip(self))] + async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { + unimplemented!() + } + + #[tracing::instrument(skip(self))] + async fn list_objects_v2( + self: Arc, + _bucket: &str, + _prefix: &str, + _continuation_token: Option, + _delimiter: Option, + _max_keys: i32, + _fetch_owner: bool, + _start_after: Option, + ) -> Result { + unimplemented!() + } + + #[tracing::instrument(skip(self))] + async fn list_object_versions( + self: Arc, + _bucket: &str, + _prefix: &str, + _marker: Option, + _version_marker: Option, + _delimiter: Option, + _max_keys: i32, + ) -> Result { + unimplemented!() + } + + #[tracing::instrument(skip(self))] + async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + self.get_disks_by_key(object).get_object_info(bucket, object, opts).await + } + #[tracing::instrument(skip(self))] async fn copy_object( &self, @@ -425,6 +459,16 @@ impl StorageAPI for Sets { ))) } + #[tracing::instrument(skip(self))] + async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { + if opts.delete_prefix && !opts.delete_prefix_object { + self.delete_prefix(bucket, object).await?; + return Ok(ObjectInfo::default()); + } + + self.get_disks_by_key(object).delete_object(bucket, object, opts).await + } + #[tracing::instrument(skip(self))] async fn delete_objects( &self, @@ -510,66 +554,22 @@ impl StorageAPI for Sets { } #[tracing::instrument(skip(self))] - async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { - if opts.delete_prefix && !opts.delete_prefix_object { - self.delete_prefix(bucket, object).await?; - return Ok(ObjectInfo::default()); - } - - self.get_disks_by_key(object).delete_object(bucket, object, opts).await - } - - #[tracing::instrument(skip(self))] - async fn list_objects_v2( - self: Arc, - _bucket: &str, - _prefix: &str, - _continuation_token: Option, - _delimiter: Option, - _max_keys: i32, - _fetch_owner: bool, - _start_after: Option, - ) -> Result { - unimplemented!() - } - - #[tracing::instrument(skip(self))] - async fn list_object_versions( - self: Arc, - _bucket: &str, - _prefix: &str, - _marker: Option, - _version_marker: Option, - _delimiter: Option, - _max_keys: i32, - ) -> Result { - unimplemented!() - } - - #[tracing::instrument(skip(self))] - async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - self.get_disks_by_key(object).get_object_info(bucket, object, opts).await - } - - #[tracing::instrument(skip(self))] - async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - self.get_disks_by_key(object).put_object_metadata(bucket, object, opts).await - } - - #[tracing::instrument(skip(self))] - async fn get_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - self.get_disks_by_key(object).get_object_tags(bucket, object, opts).await - } - #[tracing::instrument(level = "debug", skip(self))] - async fn put_object_tags(&self, bucket: &str, object: &str, tags: &str, opts: &ObjectOptions) -> Result { - self.get_disks_by_key(object) - .put_object_tags(bucket, object, tags, opts) + async fn list_multipart_uploads( + &self, + bucket: &str, + prefix: &str, + key_marker: Option, + upload_id_marker: Option, + delimiter: Option, + max_uploads: usize, + ) -> Result { + self.get_disks_by_key(prefix) + .list_multipart_uploads(bucket, prefix, key_marker, upload_id_marker, delimiter, max_uploads) .await } - #[tracing::instrument(skip(self))] - async fn delete_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - self.get_disks_by_key(object).delete_object_tags(bucket, object, opts).await + async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + self.get_disks_by_key(object).new_multipart_upload(bucket, object, opts).await } #[tracing::instrument(skip(self))] @@ -605,26 +605,6 @@ impl StorageAPI for Sets { .await } - #[tracing::instrument(skip(self))] - async fn list_multipart_uploads( - &self, - bucket: &str, - prefix: &str, - key_marker: Option, - upload_id_marker: Option, - delimiter: Option, - max_uploads: usize, - ) -> Result { - self.get_disks_by_key(prefix) - .list_multipart_uploads(bucket, prefix, key_marker, upload_id_marker, delimiter, max_uploads) - .await - } - - #[tracing::instrument(skip(self))] - async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - self.get_disks_by_key(object).new_multipart_upload(bucket, object, opts).await - } - #[tracing::instrument(skip(self))] async fn get_multipart_info( &self, @@ -670,8 +650,25 @@ impl StorageAPI for Sets { } #[tracing::instrument(skip(self))] - async fn delete_bucket(&self, _bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { - unimplemented!() + async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + self.get_disks_by_key(object).put_object_metadata(bucket, object, opts).await + } + + #[tracing::instrument(skip(self))] + async fn get_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + self.get_disks_by_key(object).get_object_tags(bucket, object, opts).await + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn put_object_tags(&self, bucket: &str, object: &str, tags: &str, opts: &ObjectOptions) -> Result { + self.get_disks_by_key(object) + .put_object_tags(bucket, object, tags, opts) + .await + } + + #[tracing::instrument(skip(self))] + async fn delete_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + self.get_disks_by_key(object).delete_object_tags(bucket, object, opts).await } #[tracing::instrument(skip(self))] diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 3c5a1eddc..5fb6553a1 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -70,7 +70,7 @@ const MAX_UPLOADS_LIST: usize = 10000; #[derive(Debug)] pub struct ECStore { - pub id: uuid::Uuid, + pub id: Uuid, // pub disks: Vec, pub disk_map: HashMap>>, pub pools: Vec>, @@ -135,7 +135,7 @@ impl ECStore { // validate_parity(partiy_count, pool_eps.drives_per_set)?; - let (disks, errs) = crate::store_init::init_disks( + let (disks, errs) = store_init::init_disks( &pool_eps.endpoints, &DiskOption { cleanup: true, @@ -169,7 +169,7 @@ impl ECStore { return Err(Error::from_string("can not get formats")); } info!("retrying get formats after {:?}", interval); - tokio::select! { + select! { _ = tokio::signal::ctrl_c() => { info!("got ctrl+c, exits"); exit(0); @@ -937,11 +937,7 @@ impl ECStore { }; if a_mod == b_mod { - if a.idx < b.idx { - return Ordering::Greater; - } else { - return Ordering::Less; - } + return if a.idx < b.idx { Ordering::Greater } else { Ordering::Less }; } b_mod.cmp(&a_mod) @@ -1190,7 +1186,7 @@ impl ObjectIO for ECStore { ) -> Result { check_get_obj_args(bucket, object)?; - let object = utils::path::encode_dir_object(object); + let object = encode_dir_object(object); if self.single_pool() { return self.pools[0].get_object_reader(bucket, object.as_str(), range, h, opts).await; @@ -1213,7 +1209,7 @@ impl ObjectIO for ECStore { async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result { check_put_object_args(bucket, object)?; - let object = utils::path::encode_dir_object(object); + let object = encode_dir_object(object); if self.single_pool() { return self.pools[0].put_object(bucket, object.as_str(), data, opts).await; @@ -1319,52 +1315,6 @@ impl StorageAPI for ECStore { madmin::StorageInfo { backend, disks } } - #[tracing::instrument(skip(self))] - async fn list_bucket(&self, opts: &BucketOptions) -> Result> { - // TODO: opts.cached - - let mut buckets = self.peer_sys.list_bucket(opts).await?; - - if !opts.no_metadata { - for bucket in buckets.iter_mut() { - if let Ok(created) = metadata_sys::created_at(&bucket.name).await { - bucket.created = Some(created); - } - } - } - Ok(buckets) - } - - #[tracing::instrument(skip(self))] - async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { - if is_meta_bucketname(bucket) { - return Err(StorageError::BucketNameInvalid(bucket.to_string()).into()); - } - - if let Err(err) = check_valid_bucket_name(bucket) { - return Err(StorageError::BucketNameInvalid(err.to_string()).into()); - } - - // TODO: nslock - - let mut opts = opts.clone(); - if !opts.force { - // FIXME: check bucket exists - opts.force = true - } - - self.peer_sys - .delete_bucket(bucket, &opts) - .await - .map_err(|e| to_object_err(e, vec![bucket]))?; - - // TODO: replication opts.srdelete_op - - // 删除 meta - self.delete_all(RUSTFS_META_BUCKET, format!("{}/{}", BUCKET_META_PREFIX, bucket).as_str()) - .await?; - Ok(()) - } #[tracing::instrument(skip(self))] async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> { if !is_meta_bucketname(bucket) { @@ -1409,6 +1359,7 @@ impl StorageAPI for ECStore { Ok(()) } + #[tracing::instrument(skip(self))] async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result { let mut info = self @@ -1425,6 +1376,101 @@ impl StorageAPI for ECStore { Ok(info) } + #[tracing::instrument(skip(self))] + async fn list_bucket(&self, opts: &BucketOptions) -> Result> { + // TODO: opts.cached + + let mut buckets = self.peer_sys.list_bucket(opts).await?; + + if !opts.no_metadata { + for bucket in buckets.iter_mut() { + if let Ok(created) = metadata_sys::created_at(&bucket.name).await { + bucket.created = Some(created); + } + } + } + Ok(buckets) + } + #[tracing::instrument(skip(self))] + async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { + if is_meta_bucketname(bucket) { + return Err(StorageError::BucketNameInvalid(bucket.to_string()).into()); + } + + if let Err(err) = check_valid_bucket_name(bucket) { + return Err(StorageError::BucketNameInvalid(err.to_string()).into()); + } + + // TODO: nslock + + let mut opts = opts.clone(); + if !opts.force { + // FIXME: check bucket exists + opts.force = true + } + + self.peer_sys + .delete_bucket(bucket, &opts) + .await + .map_err(|e| to_object_err(e, vec![bucket]))?; + + // TODO: replication opts.srdelete_op + + // 删除 meta + self.delete_all(RUSTFS_META_BUCKET, format!("{}/{}", BUCKET_META_PREFIX, bucket).as_str()) + .await?; + Ok(()) + } + + // @continuation_token marker + // @start_after as marker when continuation_token empty + // @delimiter default="/", empty when recursive + // @max_keys limit + #[tracing::instrument(skip(self))] + async fn list_objects_v2( + self: Arc, + bucket: &str, + prefix: &str, + continuation_token: Option, + delimiter: Option, + max_keys: i32, + fetch_owner: bool, + start_after: Option, + ) -> Result { + self.inner_list_objects_v2(bucket, prefix, continuation_token, delimiter, max_keys, fetch_owner, start_after) + .await + } + + #[tracing::instrument(skip(self))] + async fn list_object_versions( + self: Arc, + bucket: &str, + prefix: &str, + marker: Option, + version_marker: Option, + delimiter: Option, + max_keys: i32, + ) -> Result { + self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys) + .await + } + + #[tracing::instrument(skip(self))] + async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + check_object_args(bucket, object)?; + + let object = encode_dir_object(object); + + if self.single_pool() { + return self.pools[0].get_object_info(bucket, object.as_str(), opts).await; + } + + // TODO: nslock + + let (info, _) = self.get_latest_object_info_with_idx(bucket, object.as_str(), opts).await?; + + Ok(info) + } // TODO: review #[tracing::instrument(skip(self))] @@ -1441,8 +1487,8 @@ impl StorageAPI for ECStore { check_copy_obj_args(src_bucket, src_object)?; check_copy_obj_args(dst_bucket, dst_object)?; - let src_object = utils::path::encode_dir_object(src_object); - let dst_object = utils::path::encode_dir_object(dst_object); + let src_object = encode_dir_object(src_object); + let dst_object = encode_dir_object(dst_object); let cp_src_dst_same = path_join_buf(&[src_bucket, &src_object]) == path_join_buf(&[dst_bucket, &dst_object]); @@ -1496,7 +1542,76 @@ impl StorageAPI for ECStore { "put_object_reader is none".to_owned(), ))) } + #[tracing::instrument(skip(self))] + async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { + check_del_obj_args(bucket, object)?; + if opts.delete_prefix { + self.delete_prefix(bucket, object).await?; + return Ok(ObjectInfo::default()); + } + + // TODO: nslock + + let object = encode_dir_object(object); + let object = object.as_str(); + + // 查询在哪个 pool + let (mut pinfo, errs) = self + .get_pool_info_existing_with_opts(bucket, object, &opts) + .await + .map_err(|e| { + if is_err_read_quorum(&e) { + Error::new(StorageError::InsufficientWriteQuorum) + } else { + e + } + })?; + + if pinfo.object_info.delete_marker && opts.version_id.is_none() { + pinfo.object_info.name = decode_dir_object(object); + return Ok(pinfo.object_info); + } + + if opts.data_movement && opts.src_pool_idx == pinfo.index { + return Err(Error::new(StorageError::DataMovementOverwriteErr( + bucket.to_owned(), + object.to_owned(), + opts.version_id.unwrap_or_default(), + ))); + } + + if opts.data_movement { + let mut obj = self.pools[pinfo.index].delete_object(bucket, object, opts).await?; + obj.name = decode_dir_object(obj.name.as_str()); + return Ok(obj); + } + + if !errs.is_empty() && !opts.versioned && !opts.version_suspended { + return self.delete_object_from_all_pools(bucket, object, &opts, errs).await; + } + + for pool in self.pools.iter() { + match pool.delete_object(bucket, object, opts.clone()).await { + Ok(res) => { + let mut obj = res; + obj.name = decode_dir_object(object); + return Ok(obj); + } + Err(err) => { + if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) { + return Err(err); + } + } + } + } + + if let Some(ver) = opts.version_id { + return Err(Error::new(StorageError::VersionNotFound(bucket.to_owned(), object.to_owned(), ver))); + } + + Err(Error::new(StorageError::ObjectNotFound(bucket.to_owned(), object.to_owned()))) + } // TODO: review #[tracing::instrument(skip(self))] async fn delete_objects( @@ -1510,7 +1625,7 @@ impl StorageAPI for ECStore { .iter() .map(|v| { let mut v = v.clone(); - v.object_name = utils::path::encode_dir_object(v.object_name.as_str()); + v.object_name = encode_dir_object(v.object_name.as_str()); v }) .collect(); @@ -1580,7 +1695,7 @@ impl StorageAPI for ECStore { del_objects[i] = DeletedObject { delete_marker: pinfo.object_info.delete_marker, delete_marker_version_id: pinfo.object_info.version_id.map(|v| v.to_string()), - object_name: utils::path::decode_dir_object(&pinfo.object_info.name), + object_name: decode_dir_object(&pinfo.object_info.name), delete_marker_mtime: pinfo.object_info.mod_time, ..Default::default() }; @@ -1607,7 +1722,7 @@ impl StorageAPI for ECStore { if let Some(obj) = objects.get(i) { del_objects[i] = DeletedObject { - object_name: utils::path::decode_dir_object(&obj.object_name), + object_name: decode_dir_object(&obj.object_name), version_id: obj.version_id.map(|v| v.to_string()), ..Default::default() } @@ -1641,7 +1756,7 @@ impl StorageAPI for ECStore { } let mut dobj = pdel_objs.get(i).unwrap().clone(); - dobj.object_name = utils::path::decode_dir_object(&dobj.object_name); + dobj.object_name = decode_dir_object(&dobj.object_name); del_objects[obj_idx] = dobj; } @@ -1652,246 +1767,6 @@ impl StorageAPI for ECStore { Ok((del_objects, del_errs)) } - #[tracing::instrument(skip(self))] - async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { - check_del_obj_args(bucket, object)?; - - if opts.delete_prefix { - self.delete_prefix(bucket, object).await?; - return Ok(ObjectInfo::default()); - } - - // TODO: nslock - - let object = utils::path::encode_dir_object(object); - let object = object.as_str(); - - // 查询在哪个 pool - let (mut pinfo, errs) = self - .get_pool_info_existing_with_opts(bucket, object, &opts) - .await - .map_err(|e| { - if is_err_read_quorum(&e) { - Error::new(StorageError::InsufficientWriteQuorum) - } else { - e - } - })?; - - if pinfo.object_info.delete_marker && opts.version_id.is_none() { - pinfo.object_info.name = utils::path::decode_dir_object(object); - return Ok(pinfo.object_info); - } - - if opts.data_movement && opts.src_pool_idx == pinfo.index { - return Err(Error::new(StorageError::DataMovementOverwriteErr( - bucket.to_owned(), - object.to_owned(), - opts.version_id.unwrap_or_default(), - ))); - } - - if opts.data_movement { - let mut obj = self.pools[pinfo.index].delete_object(bucket, object, opts).await?; - obj.name = decode_dir_object(obj.name.as_str()); - return Ok(obj); - } - - if !errs.is_empty() && !opts.versioned && !opts.version_suspended { - return self.delete_object_from_all_pools(bucket, object, &opts, errs).await; - } - - for pool in self.pools.iter() { - match pool.delete_object(bucket, object, opts.clone()).await { - Ok(res) => { - let mut obj = res; - obj.name = utils::path::decode_dir_object(object); - return Ok(obj); - } - Err(err) => { - if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) { - return Err(err); - } - } - } - } - - if let Some(ver) = opts.version_id { - return Err(Error::new(StorageError::VersionNotFound(bucket.to_owned(), object.to_owned(), ver))); - } - - Err(Error::new(StorageError::ObjectNotFound(bucket.to_owned(), object.to_owned()))) - } - - // @continuation_token marker - // @start_after as marker when continuation_token empty - // @delimiter default="/", empty when recursive - // @max_keys limit - #[tracing::instrument(skip(self))] - async fn list_objects_v2( - self: Arc, - bucket: &str, - prefix: &str, - continuation_token: Option, - delimiter: Option, - max_keys: i32, - fetch_owner: bool, - start_after: Option, - ) -> Result { - self.inner_list_objects_v2(bucket, prefix, continuation_token, delimiter, max_keys, fetch_owner, start_after) - .await - } - #[tracing::instrument(skip(self))] - async fn list_object_versions( - self: Arc, - bucket: &str, - prefix: &str, - marker: Option, - version_marker: Option, - delimiter: Option, - max_keys: i32, - ) -> Result { - self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys) - .await - } - #[tracing::instrument(skip(self))] - async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - check_object_args(bucket, object)?; - - let object = utils::path::encode_dir_object(object); - - if self.single_pool() { - return self.pools[0].get_object_info(bucket, object.as_str(), opts).await; - } - - // TODO: nslock - - let (info, _) = self.get_latest_object_info_with_idx(bucket, object.as_str(), opts).await?; - - Ok(info) - } - - #[tracing::instrument(skip(self))] - async fn get_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - let object = encode_dir_object(object); - - if self.single_pool() { - return self.pools[0].get_object_tags(bucket, object.as_str(), opts).await; - } - - let (oi, _) = self.get_latest_object_info_with_idx(bucket, &object, opts).await?; - - Ok(oi.user_tags) - } - - #[tracing::instrument(skip(self))] - async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - let object = encode_dir_object(object); - if self.single_pool() { - return self.pools[0].put_object_metadata(bucket, object.as_str(), opts).await; - } - - let mut opts = opts.clone(); - opts.metadata_chg = true; - - let idx = self.get_pool_idx_existing_with_opts(bucket, object.as_str(), &opts).await?; - - self.pools[idx].put_object_metadata(bucket, object.as_str(), &opts).await - } - - #[tracing::instrument(level = "debug", skip(self))] - async fn put_object_tags(&self, bucket: &str, object: &str, tags: &str, opts: &ObjectOptions) -> Result { - let object = encode_dir_object(object); - - if self.single_pool() { - return self.pools[0].put_object_tags(bucket, object.as_str(), tags, opts).await; - } - - let idx = self.get_pool_idx_existing_with_opts(bucket, object.as_str(), opts).await?; - - self.pools[idx].put_object_tags(bucket, object.as_str(), tags, opts).await - } - #[tracing::instrument(skip(self))] - async fn delete_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - let object = encode_dir_object(object); - - if self.single_pool() { - return self.pools[0].delete_object_tags(bucket, object.as_str(), opts).await; - } - - let idx = self.get_pool_idx_existing_with_opts(bucket, object.as_str(), opts).await?; - - self.pools[idx].delete_object_tags(bucket, object.as_str(), opts).await - } - - #[tracing::instrument(skip(self))] - async fn copy_object_part( - &self, - src_bucket: &str, - src_object: &str, - _dst_bucket: &str, - _dst_object: &str, - _upload_id: &str, - _part_id: usize, - _start_offset: i64, - _length: i64, - _src_info: &ObjectInfo, - _src_opts: &ObjectOptions, - _dst_opts: &ObjectOptions, - ) -> Result<()> { - check_new_multipart_args(src_bucket, src_object)?; - - // TODO: PutObjectReader - // self.put_object_part(dst_bucket, dst_object, upload_id, part_id, data, opts) - - unimplemented!() - } - #[tracing::instrument(skip(self, data))] - async fn put_object_part( - &self, - bucket: &str, - object: &str, - upload_id: &str, - part_id: usize, - data: &mut PutObjReader, - opts: &ObjectOptions, - ) -> Result { - check_put_object_part_args(bucket, object, upload_id)?; - - if self.single_pool() { - return self.pools[0] - .put_object_part(bucket, object, upload_id, part_id, data, opts) - .await; - } - - for pool in self.pools.iter() { - if self.is_suspended(pool.pool_idx).await { - continue; - } - let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await { - Ok(res) => return Ok(res), - Err(err) => { - if is_err_invalid_upload_id(&err) { - None - } else { - Some(err) - } - } - }; - - if let Some(err) = err { - error!("put_object_part err: {:?}", err); - return Err(err); - } - } - - Err(Error::new(StorageError::InvalidUploadID( - bucket.to_owned(), - object.to_owned(), - upload_id.to_owned(), - ))) - } - #[tracing::instrument(skip(self))] async fn list_multipart_uploads( &self, @@ -1977,6 +1852,74 @@ impl StorageAPI for ECStore { self.pools[idx].new_multipart_upload(bucket, object, opts).await } + #[tracing::instrument(skip(self))] + async fn copy_object_part( + &self, + src_bucket: &str, + src_object: &str, + _dst_bucket: &str, + _dst_object: &str, + _upload_id: &str, + _part_id: usize, + _start_offset: i64, + _length: i64, + _src_info: &ObjectInfo, + _src_opts: &ObjectOptions, + _dst_opts: &ObjectOptions, + ) -> Result<()> { + check_new_multipart_args(src_bucket, src_object)?; + + // TODO: PutObjectReader + // self.put_object_part(dst_bucket, dst_object, upload_id, part_id, data, opts) + + unimplemented!() + } + #[tracing::instrument(skip(self, data))] + async fn put_object_part( + &self, + bucket: &str, + object: &str, + upload_id: &str, + part_id: usize, + data: &mut PutObjReader, + opts: &ObjectOptions, + ) -> Result { + check_put_object_part_args(bucket, object, upload_id)?; + + if self.single_pool() { + return self.pools[0] + .put_object_part(bucket, object, upload_id, part_id, data, opts) + .await; + } + + for pool in self.pools.iter() { + if self.is_suspended(pool.pool_idx).await { + continue; + } + let err = match pool.put_object_part(bucket, object, upload_id, part_id, data, opts).await { + Ok(res) => return Ok(res), + Err(err) => { + if is_err_invalid_upload_id(&err) { + None + } else { + Some(err) + } + } + }; + + if let Some(err) = err { + error!("put_object_part err: {:?}", err); + return Err(err); + } + } + + Err(Error::new(StorageError::InvalidUploadID( + bucket.to_owned(), + object.to_owned(), + upload_id.to_owned(), + ))) + } + #[tracing::instrument(skip(self))] async fn get_multipart_info( &self, @@ -1995,16 +1938,16 @@ impl StorageAPI for ECStore { continue; } - match pool.get_multipart_info(bucket, object, upload_id, opts).await { - Ok(res) => return Ok(res), + return match pool.get_multipart_info(bucket, object, upload_id, opts).await { + Ok(res) => Ok(res), Err(err) => { if is_err_invalid_upload_id(&err) { continue; } - return Err(err); + Err(err) } - } + }; } Err(Error::new(StorageError::InvalidUploadID( @@ -2051,6 +1994,7 @@ impl StorageAPI for ECStore { upload_id.to_owned(), ))) } + #[tracing::instrument(skip(self))] async fn complete_multipart_upload( &self, @@ -2118,6 +2062,58 @@ impl StorageAPI for ECStore { } counts } + #[tracing::instrument(skip(self))] + async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + let object = encode_dir_object(object); + if self.single_pool() { + return self.pools[0].put_object_metadata(bucket, object.as_str(), opts).await; + } + + let mut opts = opts.clone(); + opts.metadata_chg = true; + + let idx = self.get_pool_idx_existing_with_opts(bucket, object.as_str(), &opts).await?; + + self.pools[idx].put_object_metadata(bucket, object.as_str(), &opts).await + } + #[tracing::instrument(skip(self))] + async fn get_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + let object = encode_dir_object(object); + + if self.single_pool() { + return self.pools[0].get_object_tags(bucket, object.as_str(), opts).await; + } + + let (oi, _) = self.get_latest_object_info_with_idx(bucket, &object, opts).await?; + + Ok(oi.user_tags) + } + + #[tracing::instrument(level = "debug", skip(self))] + async fn put_object_tags(&self, bucket: &str, object: &str, tags: &str, opts: &ObjectOptions) -> Result { + let object = encode_dir_object(object); + + if self.single_pool() { + return self.pools[0].put_object_tags(bucket, object.as_str(), tags, opts).await; + } + + let idx = self.get_pool_idx_existing_with_opts(bucket, object.as_str(), opts).await?; + + self.pools[idx].put_object_tags(bucket, object.as_str(), tags, opts).await + } + + #[tracing::instrument(skip(self))] + async fn delete_object_tags(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + let object = encode_dir_object(object); + + if self.single_pool() { + return self.pools[0].delete_object_tags(bucket, object.as_str(), opts).await; + } + + let idx = self.get_pool_idx_existing_with_opts(bucket, object.as_str(), opts).await?; + + self.pools[idx].delete_object_tags(bucket, object.as_str(), opts).await + } #[tracing::instrument(skip(self))] async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { @@ -2167,7 +2163,7 @@ impl StorageAPI for ECStore { opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { info!("ECStore heal_object"); - let object = utils::path::encode_dir_object(object); + let object = encode_dir_object(object); let mut futures = Vec::with_capacity(self.pools.len()); for pool in self.pools.iter() { @@ -2195,7 +2191,7 @@ impl StorageAPI for ECStore { match res { Ok((result, err)) => { let mut result = result; - result.object = utils::path::decode_dir_object(&result.object); + result.object = decode_dir_object(&result.object); ress.push(result); errs.push(err); } @@ -2265,10 +2261,10 @@ impl StorageAPI for ECStore { let fivs = match entry.file_info_versions(&bucket) { Ok(fivs) => fivs, Err(_) => { - if is_meta { - return HealSequence::heal_meta_object(hs_clone.clone(), &bucket, &entry.name, "", scan_mode).await; + return if is_meta { + HealSequence::heal_meta_object(hs_clone.clone(), &bucket, &entry.name, "", scan_mode).await } else { - return HealSequence::heal_object(hs_clone.clone(), &bucket, &entry.name, "", scan_mode).await; + HealSequence::heal_object(hs_clone.clone(), &bucket, &entry.name, "", scan_mode).await } } }; @@ -2362,7 +2358,7 @@ impl StorageAPI for ECStore { #[tracing::instrument(skip(self))] async fn check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> { - let object = utils::path::encode_dir_object(object); + let object = encode_dir_object(object); if self.single_pool() { return self.pools[0].check_abandoned_parts(bucket, &object, opts).await; } diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 598dc8818..ca84354b9 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -375,7 +375,7 @@ pub struct ChecksumInfo { pub const DEFAULT_BITROT_ALGO: BitrotAlgorithm = BitrotAlgorithm::HighwayHash256S; #[derive(Serialize, Deserialize, Debug, PartialEq, Default, Clone, Eq, Hash)] -// BitrotAlgorithm specifies a algorithm used for bitrot protection. +// BitrotAlgorithm specifies an algorithm used for bitrot protection. pub enum BitrotAlgorithm { // SHA256 represents the SHA-256 hash function SHA256, @@ -465,23 +465,23 @@ impl GetObjectReader { if let Some(rs) = rs { let (off, length) = rs.get_offset_length(oi.size)?; - return Ok(( + Ok(( GetObjectReader { stream: reader, object_info: oi.clone(), }, off, length, - )); + )) } else { - return Ok(( + Ok(( GetObjectReader { stream: reader, object_info: oi.clone(), }, 0, oi.size, - )); + )) } } pub async fn read_all(&mut self) -> Result> { diff --git a/ecstore/src/store_list_objects.rs b/ecstore/src/store_list_objects.rs index 445c56f0c..cc35e213e 100644 --- a/ecstore/src/store_list_objects.rs +++ b/ecstore/src/store_list_objects.rs @@ -277,13 +277,13 @@ impl ECStore { }; }; - let mut list_result = match self.list_path(&opts).await { - Ok(res) => res, - Err(err) => MetaCacheEntriesSortedResult { + let mut list_result = self + .list_path(&opts) + .await + .unwrap_or_else(|err| MetaCacheEntriesSortedResult { err: Some(err), ..Default::default() - }, - }; + }); if let Some(err) = &list_result.err { if !is_err_eof(err) { diff --git a/iam/src/error.rs b/iam/src/error.rs index c0314bf9f..4f42a0843 100644 --- a/iam/src/error.rs +++ b/iam/src/error.rs @@ -7,7 +7,7 @@ pub enum Error { #[error(transparent)] PolicyError(#[from] PolicyError), - #[error("ecsotre error: {0}")] + #[error("ecstore error: {0}")] EcstoreError(common::error::Error), #[error("{0}")] diff --git a/iam/src/sys.rs b/iam/src/sys.rs index 175dae6bc..0fe10346b 100644 --- a/iam/src/sys.rs +++ b/iam/src/sys.rs @@ -226,13 +226,13 @@ impl IamSys { }; let mut m: HashMap = HashMap::new(); - m.insert("parent".to_owned(), serde_json::Value::String(parent_user.to_owned())); + m.insert("parent".to_owned(), Value::String(parent_user.to_owned())); if !policy_buf.is_empty() { - m.insert(SESSION_POLICY_NAME.to_owned(), serde_json::Value::String(base64_encode(&policy_buf))); - m.insert(iam_policy_claim_name_sa(), serde_json::Value::String(EMBEDDED_POLICY_TYPE.to_owned())); + m.insert(SESSION_POLICY_NAME.to_owned(), Value::String(base64_encode(&policy_buf))); + m.insert(iam_policy_claim_name_sa(), Value::String(EMBEDDED_POLICY_TYPE.to_owned())); } else { - m.insert(iam_policy_claim_name_sa(), serde_json::Value::String(INHERITED_POLICY_TYPE.to_owned())); + m.insert(iam_policy_claim_name_sa(), Value::String(INHERITED_POLICY_TYPE.to_owned())); } if let Some(claims) = opts.claims { @@ -246,7 +246,7 @@ impl IamSys { // set expiration time default to 1 hour m.insert( "exp".to_string(), - serde_json::Value::Number(serde_json::Number::from( + Value::Number(serde_json::Number::from( opts.expiration .map_or(OffsetDateTime::now_utc().unix_timestamp() + 3600, |t| t.unix_timestamp()), )), diff --git a/iam/src/utils.rs b/iam/src/utils.rs index 4890e55b7..90a82e819 100644 --- a/iam/src/utils.rs +++ b/iam/src/utils.rs @@ -23,7 +23,7 @@ pub fn gen_access_key(length: usize) -> Result { Ok(result) } -pub fn gen_secret_key(length: usize) -> crate::Result { +pub fn gen_secret_key(length: usize) -> Result { use base64_simd::URL_SAFE_NO_PAD; if length < 8 { diff --git a/policy/src/arn.rs b/policy/src/arn.rs index 2337208ac..472ca84f9 100644 --- a/policy/src/arn.rs +++ b/policy/src/arn.rs @@ -17,7 +17,7 @@ pub struct ARN { impl ARN { pub fn new_iam_role_arn(resource_id: &str, server_region: &str) -> Result { - let valid_resource_id_regex = Regex::new(r"^[A-Za-z0-9_/\.-]+$").unwrap(); + let valid_resource_id_regex = Regex::new(r"^[A-Za-z0-9_/\.-]+$")?; if !valid_resource_id_regex.is_match(resource_id) { return Err(Error::msg("ARN resource ID invalid")); } @@ -57,7 +57,7 @@ impl ARN { return Err(Error::msg("ARN resource type invalid")); } - let valid_resource_id_regex = Regex::new(r"^[A-Za-z0-9_/\.-]+$").unwrap(); + let valid_resource_id_regex = Regex::new(r"^[A-Za-z0-9_/\.-]+$")?; if !valid_resource_id_regex.is_match(res[1]) { return Err(Error::msg("ARN resource ID invalid")); } diff --git a/policy/src/policy/function/date.rs b/policy/src/policy/function/date.rs index e3abaf880..4f02fb891 100644 --- a/policy/src/policy/function/date.rs +++ b/policy/src/policy/function/date.rs @@ -39,7 +39,7 @@ impl Serialize for DateFuncValue { &self .0 .format(&Rfc3339) - .map_err(|e| S::Error::custom(format!("format datetime failed: {e:?}")))?, + .map_err(|e| Error::custom(format!("format datetime failed: {e:?}")))?, ) } } diff --git a/policy/src/policy/function/key_name.rs b/policy/src/policy/function/key_name.rs index 4457cee08..57cbae43d 100644 --- a/policy/src/policy/function/key_name.rs +++ b/policy/src/policy/function/key_name.rs @@ -309,7 +309,6 @@ pub enum AwsKeyName { #[cfg(test)] mod tests { use super::*; - use crate::policy::Error; use serde::Deserialize; use test_case::test_case; @@ -330,7 +329,7 @@ mod tests { #[test_case("ldap:us")] #[test_case("DurationSeconds")] fn key_name_from_str_failed(val: &str) { - assert_eq!(KeyName::try_from(val), Err(Error::InvalidKeyName(val.to_string()))); + assert_eq!(KeyName::try_from(val), Err(InvalidKeyName(val.to_string()))); } #[test_case("s3:x-amz-copy-source", KeyName::S3(S3KeyName::S3XAmzCopySource))] diff --git a/rustfs/src/console.rs b/rustfs/src/console.rs index 0523991d6..8d939523b 100644 --- a/rustfs/src/console.rs +++ b/rustfs/src/console.rs @@ -174,9 +174,9 @@ async fn license_handler() -> impl IntoResponse { .unwrap() } -fn _is_private_ip(ip: std::net::IpAddr) -> bool { +fn _is_private_ip(ip: IpAddr) -> bool { match ip { - std::net::IpAddr::V4(ip) => { + IpAddr::V4(ip) => { let octets = ip.octets(); // 10.0.0.0/8 octets[0] == 10 || @@ -185,7 +185,7 @@ fn _is_private_ip(ip: std::net::IpAddr) -> bool { // 192.168.0.0/16 (octets[0] == 192 && octets[1] == 168) } - std::net::IpAddr::V6(_) => false, + IpAddr::V6(_) => false, } } diff --git a/rustfs/src/license.rs b/rustfs/src/license.rs index 2206c40e2..d805dc159 100644 --- a/rustfs/src/license.rs +++ b/rustfs/src/license.rs @@ -34,7 +34,7 @@ pub fn get_license() -> Option { #[allow(unreachable_code)] pub fn license_check() -> Result<()> { return Ok(()); - let inval_license = LICENSE.get().map(|token| { + let invalid_license = LICENSE.get().map(|token| { if token.expired < SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs() { error!("License expired"); return Err(Error::from_string("Incorrect license, please contact RustFS.".to_string())); @@ -43,7 +43,7 @@ pub fn license_check() -> Result<()> { Ok(()) }); - // let inval_license = config::get_config().license.as_ref().map(|license| { + // let invalid_license = config::get_config().license.as_ref().map(|license| { // if license.is_empty() { // error!("License is empty"); // return Err(Error::from_string("Incorrect license, please contact RustFS.".to_string())); @@ -58,7 +58,7 @@ pub fn license_check() -> Result<()> { // Ok(()) // }); - if inval_license.is_none() || inval_license.is_some_and(|v| v.is_err()) { + if invalid_license.is_none() || invalid_license.is_some_and(|v| v.is_err()) { return Err(Error::from_string("Incorrect license, please contact RustFS.".to_string())); } diff --git a/s3select/api/src/object_store.rs b/s3select/api/src/object_store.rs index 52132bc1a..e70a8b70d 100644 --- a/s3select/api/src/object_store.rs +++ b/s3select/api/src/object_store.rs @@ -203,7 +203,7 @@ impl AsyncRead for ConvertStream { self: Pin<&mut Self>, cx: &mut std::task::Context<'_>, buf: &mut tokio::io::ReadBuf<'_>, - ) -> std::task::Poll> { + ) -> Poll> { let me = self.project(); ready!(Pin::new(&mut *me.inner).poll_read(cx, buf))?; let bytes = buf.filled(); diff --git a/s3select/api/src/query/execution.rs b/s3select/api/src/query/execution.rs index 10c48acc3..99fb9671f 100644 --- a/s3select/api/src/query/execution.rs +++ b/s3select/api/src/query/execution.rs @@ -92,7 +92,7 @@ impl Output { } impl Stream for Output { - type Item = std::result::Result; + type Item = Result; fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { let this = self.get_mut();