diff --git a/DEVELOPMENT.md b/DEVELOPMENT.md index 76884ce26..be58a46e2 100644 --- a/DEVELOPMENT.md +++ b/DEVELOPMENT.md @@ -11,21 +11,25 @@ Before every commit, you **MUST**: 1. **Format your code**: + ```bash cargo fmt --all ``` 2. **Verify formatting**: + ```bash cargo fmt --all --check ``` 3. **Pass clippy checks**: + ```bash cargo clippy --all-targets --all-features -- -D warnings ``` 4. **Ensure compilation**: + ```bash cargo check --all-targets ``` @@ -136,6 +140,7 @@ Install the `rust-analyzer` extension and add to your `settings.json`: #### Other IDEs Configure your IDE to: + - Use the project's `rustfmt.toml` configuration - Format on save - Run clippy checks diff --git a/crates/common/src/last_minute.rs b/crates/common/src/last_minute.rs index 7ab189bc7..b5b22eb3a 100644 --- a/crates/common/src/last_minute.rs +++ b/crates/common/src/last_minute.rs @@ -568,11 +568,14 @@ mod tests { let mut latency = LastMinuteLatency::default(); // Add data at time 1000 - latency.add_all(1000, &AccElem { - total: 10, - size: 0, - n: 1, - }); + latency.add_all( + 1000, + &AccElem { + total: 10, + size: 0, + n: 1, + }, + ); // Forward to time 1030 (30 seconds later) latency.forward_to(1030); diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 8a86f3654..791134da9 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -236,10 +236,13 @@ impl BucketMetadataSys { futures.push(async move { sleep(Duration::from_millis(30)).await; let _ = api - .heal_bucket(&bucket, &HealOpts { - recreate: true, - ..Default::default() - }) + .heal_bucket( + &bucket, + &HealOpts { + recreate: true, + ..Default::default() + }, + ) .await; load_bucket_metadata(self.api.clone(), bucket.as_str()).await }); diff --git a/crates/ecstore/src/client/api_bucket_policy.rs b/crates/ecstore/src/client/api_bucket_policy.rs index 7b9dec295..8ed9c6065 100644 --- a/crates/ecstore/src/client/api_bucket_policy.rs +++ b/crates/ecstore/src/client/api_bucket_policy.rs @@ -74,23 +74,26 @@ impl TransitionClient { url_values.insert("policy".to_string(), "".to_string()); let resp = self - .execute_method(http::Method::DELETE, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - query_values: url_values, - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - object_name: "".to_string(), - custom_header: HeaderMap::new(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::DELETE, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + query_values: url_values, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + object_name: "".to_string(), + custom_header: HeaderMap::new(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; //defer closeResponse(resp) @@ -111,23 +114,26 @@ impl TransitionClient { url_values.insert("policy".to_string(), "".to_string()); let resp = self - .execute_method(http::Method::GET, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - query_values: url_values, - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - object_name: "".to_string(), - custom_header: HeaderMap::new(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::GET, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + query_values: url_values, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + object_name: "".to_string(), + custom_header: HeaderMap::new(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; let policy = String::from_utf8_lossy(&resp.body().bytes().expect("err").to_vec()).to_string(); diff --git a/crates/ecstore/src/client/api_get_object.rs b/crates/ecstore/src/client/api_get_object.rs index cce4b1718..80759b2cc 100644 --- a/crates/ecstore/src/client/api_get_object.rs +++ b/crates/ecstore/src/client/api_get_object.rs @@ -43,23 +43,26 @@ impl TransitionClient { opts: &GetObjectOptions, ) -> Result<(ObjectInfo, HeaderMap, ReadCloser), std::io::Error> { let resp = self - .execute_method(http::Method::GET, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - query_values: opts.to_query_values(), - custom_header: opts.header(), - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::GET, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + query_values: opts.to_query_values(), + custom_header: opts.header(), + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; let resp = &resp; diff --git a/crates/ecstore/src/client/api_get_object_acl.rs b/crates/ecstore/src/client/api_get_object_acl.rs index 898e8aab5..0ee14a426 100644 --- a/crates/ecstore/src/client/api_get_object_acl.rs +++ b/crates/ecstore/src/client/api_get_object_acl.rs @@ -63,23 +63,26 @@ impl TransitionClient { let mut url_values = HashMap::new(); url_values.insert("acl".to_string(), "".to_string()); let mut resp = self - .execute_method(http::Method::GET, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - query_values: url_values, - custom_header: HeaderMap::new(), - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::GET, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + query_values: url_values, + custom_header: HeaderMap::new(), + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; if resp.status() != http::StatusCode::OK { diff --git a/crates/ecstore/src/client/api_get_object_attributes.rs b/crates/ecstore/src/client/api_get_object_attributes.rs index 15c25f347..f236118d6 100644 --- a/crates/ecstore/src/client/api_get_object_attributes.rs +++ b/crates/ecstore/src/client/api_get_object_attributes.rs @@ -193,23 +193,26 @@ impl TransitionClient { }*/ let mut resp = self - .execute_method(http::Method::HEAD, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - query_values: url_values, - custom_header: headers, - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - content_md5_base64: "".to_string(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::HEAD, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + query_values: url_values, + custom_header: headers, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + content_md5_base64: "".to_string(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; let h = resp.headers(); diff --git a/crates/ecstore/src/client/api_list.rs b/crates/ecstore/src/client/api_list.rs index d2618c4a6..f978f0634 100644 --- a/crates/ecstore/src/client/api_list.rs +++ b/crates/ecstore/src/client/api_list.rs @@ -76,23 +76,26 @@ impl TransitionClient { } let mut resp = self - .execute_method(http::Method::GET, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: "".to_string(), - query_values: url_values, - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - custom_header: headers, - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::GET, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: "".to_string(), + query_values: url_values, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + custom_header: headers, + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; if resp.status() != StatusCode::OK { return Err(std::io::Error::other(http_resp_to_error_response(resp, vec![], bucket_name, ""))); diff --git a/crates/ecstore/src/client/api_remove.rs b/crates/ecstore/src/client/api_remove.rs index 14f6cd2a7..a6845229a 100644 --- a/crates/ecstore/src/client/api_remove.rs +++ b/crates/ecstore/src/client/api_remove.rs @@ -78,23 +78,26 @@ impl TransitionClient { let headers = HeaderMap::new(); let resp = self - .execute_method(Method::DELETE, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - custom_header: headers, - object_name: "".to_string(), - query_values: Default::default(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + Method::DELETE, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + custom_header: headers, + object_name: "".to_string(), + query_values: Default::default(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; { @@ -106,23 +109,26 @@ impl TransitionClient { pub async fn remove_bucket(&self, bucket_name: &str) -> Result<(), std::io::Error> { let resp = self - .execute_method(http::Method::DELETE, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - custom_header: Default::default(), - object_name: "".to_string(), - query_values: Default::default(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::DELETE, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + custom_header: Default::default(), + object_name: "".to_string(), + query_values: Default::default(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; { @@ -157,23 +163,26 @@ impl TransitionClient { } let resp = self - .execute_method(http::Method::DELETE, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - query_values: url_values, - custom_header: headers, - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::DELETE, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + query_values: url_values, + custom_header: headers, + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; Ok(RemoveObjectResult { @@ -268,11 +277,15 @@ impl TransitionClient { while let Some(object) = objects_rx.recv().await { if has_invalid_xml_char(&object.name) { let remove_result = self - .remove_object_inner(bucket_name, &object.name, RemoveObjectOptions { - version_id: object.version_id.expect("err").to_string(), - governance_bypass: opts.governance_bypass, - ..Default::default() - }) + .remove_object_inner( + bucket_name, + &object.name, + RemoveObjectOptions { + version_id: object.version_id.expect("err").to_string(), + governance_bypass: opts.governance_bypass, + ..Default::default() + }, + ) .await?; let remove_result_clone = remove_result.clone(); if !remove_result.err.is_none() { @@ -309,23 +322,26 @@ impl TransitionClient { let remove_bytes = generate_remove_multi_objects_request(&batch); let resp = self - .execute_method(http::Method::POST, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - query_values: url_values.clone(), - content_body: ReaderImpl::Body(Bytes::from(remove_bytes.clone())), - content_length: remove_bytes.len() as i64, - content_md5_base64: base64_encode(&HashAlgorithm::Md5.hash_encode(&remove_bytes).as_ref()), - content_sha256_hex: base64_encode(&HashAlgorithm::SHA256.hash_encode(&remove_bytes).as_ref()), - custom_header: headers, - object_name: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::POST, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + query_values: url_values.clone(), + content_body: ReaderImpl::Body(Bytes::from(remove_bytes.clone())), + content_length: remove_bytes.len() as i64, + content_md5_base64: base64_encode(&HashAlgorithm::Md5.hash_encode(&remove_bytes).as_ref()), + content_sha256_hex: base64_encode(&HashAlgorithm::SHA256.hash_encode(&remove_bytes).as_ref()), + custom_header: headers, + object_name: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; let body_bytes: Vec = resp.body().bytes().expect("err").to_vec(); @@ -353,23 +369,26 @@ impl TransitionClient { url_values.insert("uploadId".to_string(), upload_id.to_string()); let resp = self - .execute_method(http::Method::DELETE, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - query_values: url_values, - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - custom_header: HeaderMap::new(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - content_md5_base64: "".to_string(), - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::DELETE, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + query_values: url_values, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + custom_header: HeaderMap::new(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; //if resp.is_some() { if resp.status() != StatusCode::NO_CONTENT { diff --git a/crates/ecstore/src/client/api_restore.rs b/crates/ecstore/src/client/api_restore.rs index 48999432b..9dc5fead7 100644 --- a/crates/ecstore/src/client/api_restore.rs +++ b/crates/ecstore/src/client/api_restore.rs @@ -141,23 +141,26 @@ impl TransitionClient { let restore_request_buffer = Bytes::from(restore_request_bytes.clone()); let resp = self - .execute_method(http::Method::HEAD, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - query_values: url_values, - custom_header: HeaderMap::new(), - content_sha256_hex: "".to_string(), //sum_sha256_hex(&restore_request_bytes), - content_md5_base64: "".to_string(), //sum_md5_base64(&restore_request_bytes), - content_body: ReaderImpl::Body(restore_request_buffer), - content_length: restore_request_bytes.len() as i64, - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::HEAD, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + query_values: url_values, + custom_header: HeaderMap::new(), + content_sha256_hex: "".to_string(), //sum_sha256_hex(&restore_request_bytes), + content_md5_base64: "".to_string(), //sum_md5_base64(&restore_request_bytes), + content_body: ReaderImpl::Body(restore_request_buffer), + content_length: restore_request_bytes.len() as i64, + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await?; let b = resp.body().bytes().expect("err").to_vec(); diff --git a/crates/ecstore/src/client/api_stat.rs b/crates/ecstore/src/client/api_stat.rs index 8e0f62f50..99eed21f0 100644 --- a/crates/ecstore/src/client/api_stat.rs +++ b/crates/ecstore/src/client/api_stat.rs @@ -35,23 +35,26 @@ use s3s::header::{X_AMZ_DELETE_MARKER, X_AMZ_VERSION_ID}; impl TransitionClient { pub async fn bucket_exists(&self, bucket_name: &str) -> Result { let resp = self - .execute_method(http::Method::HEAD, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: "".to_string(), - query_values: HashMap::new(), - custom_header: HeaderMap::new(), - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - content_md5_base64: "".to_string(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::HEAD, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: "".to_string(), + query_values: HashMap::new(), + custom_header: HeaderMap::new(), + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + content_md5_base64: "".to_string(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await; if let Ok(resp) = resp { @@ -82,23 +85,26 @@ impl TransitionClient { } let resp = self - .execute_method(http::Method::HEAD, &mut RequestMetadata { - bucket_name: bucket_name.to_string(), - object_name: object_name.to_string(), - query_values: opts.to_query_values(), - custom_header: headers, - content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), - content_md5_base64: "".to_string(), - content_body: ReaderImpl::Body(Bytes::new()), - content_length: 0, - stream_sha256: false, - trailer: HeaderMap::new(), - pre_sign_url: Default::default(), - add_crc: Default::default(), - extra_pre_sign_header: Default::default(), - bucket_location: Default::default(), - expires: Default::default(), - }) + .execute_method( + http::Method::HEAD, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: object_name.to_string(), + query_values: opts.to_query_values(), + custom_header: headers, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + content_md5_base64: "".to_string(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) .await; match resp { diff --git a/crates/ecstore/src/client/object_handlers_common.rs b/crates/ecstore/src/client/object_handlers_common.rs index 75894e437..6c380aab5 100644 --- a/crates/ecstore/src/client/object_handlers_common.rs +++ b/crates/ecstore/src/client/object_handlers_common.rs @@ -31,10 +31,14 @@ pub async fn delete_object_versions(api: ECStore, bucket: &str, to_del: &[Object remaining = &[]; } let vc = BucketVersioningSys::get(bucket).await.expect("err!"); - let _deleted_objs = api.delete_objects(bucket, to_del.to_vec(), ObjectOptions { - //prefix_enabled_fn: vc.prefix_enabled(""), - version_suspended: vc.suspended(), - ..Default::default() - }); + let _deleted_objs = api.delete_objects( + bucket, + to_del.to_vec(), + ObjectOptions { + //prefix_enabled_fn: vc.prefix_enabled(""), + version_suspended: vc.suspended(), + ..Default::default() + }, + ); } } diff --git a/crates/ecstore/src/config/com.rs b/crates/ecstore/src/config/com.rs index b73d6e76f..11b45cbca 100644 --- a/crates/ecstore/src/config/com.rs +++ b/crates/ecstore/src/config/com.rs @@ -70,20 +70,29 @@ pub async fn read_config_with_metadata( } pub async fn save_config(api: Arc, file: &str, data: Vec) -> Result<()> { - save_config_with_opts(api, file, data, &ObjectOptions { - max_parity: true, - ..Default::default() - }) + save_config_with_opts( + api, + file, + data, + &ObjectOptions { + max_parity: true, + ..Default::default() + }, + ) .await } pub async fn delete_config(api: Arc, file: &str) -> Result<()> { match api - .delete_object(RUSTFS_META_BUCKET, file, ObjectOptions { - delete_prefix: true, - delete_prefix_object: true, - ..Default::default() - }) + .delete_object( + RUSTFS_META_BUCKET, + file, + ObjectOptions { + delete_prefix: true, + delete_prefix_object: true, + ..Default::default() + }, + ) .await { Ok(_) => Ok(()), diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index a320af419..acb3b97a7 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -2026,11 +2026,15 @@ impl DiskAPI for LocalDisk { ) -> Result<()> { if path.starts_with(SLASH_SEPARATOR) { return self - .delete(volume, path, DeleteOptions { - recursive: false, - immediate: false, - ..Default::default() - }) + .delete( + volume, + path, + DeleteOptions { + recursive: false, + immediate: false, + ..Default::default() + }, + ) .await; } diff --git a/crates/ecstore/src/erasure_coding/bitrot.rs b/crates/ecstore/src/erasure_coding/bitrot.rs index cc468c3d3..587b127c8 100644 --- a/crates/ecstore/src/erasure_coding/bitrot.rs +++ b/crates/ecstore/src/erasure_coding/bitrot.rs @@ -317,10 +317,13 @@ enum WriterType { impl std::fmt::Debug for BitrotWriterWrapper { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("BitrotWriterWrapper") - .field("writer_type", &match self.writer_type { - WriterType::InlineBuffer => "InlineBuffer", - WriterType::Other => "Other", - }) + .field( + "writer_type", + &match self.writer_type { + WriterType::InlineBuffer => "InlineBuffer", + WriterType::Other => "Other", + }, + ) .finish() } } diff --git a/crates/ecstore/src/heal/data_usage.rs b/crates/ecstore/src/heal/data_usage.rs index caf50547d..a6d569971 100644 --- a/crates/ecstore/src/heal/data_usage.rs +++ b/crates/ecstore/src/heal/data_usage.rs @@ -174,10 +174,13 @@ pub async fn load_data_usage_from_backend(store: Arc) -> Result) -> Result) { for (tier, st) in &self.tiers { - stats.insert(tier.clone(), TierStats { - total_size: st.total_size, - num_versions: st.num_versions, - num_objects: st.num_objects, - }); + stats.insert( + tier.clone(), + TierStats { + total_size: st.total_size, + num_versions: st.num_versions, + num_objects: st.num_objects, + }, + ); } } } @@ -443,10 +446,16 @@ impl DataUsageCache { let path = Path::new(BUCKET_META_PREFIX).join(name); // warn!("Loading data usage cache from backend: {}", path.display()); match store - .get_object_reader(RUSTFS_META_BUCKET, path.to_str().unwrap(), None, HeaderMap::new(), &ObjectOptions { - no_lock: true, - ..Default::default() - }) + .get_object_reader( + RUSTFS_META_BUCKET, + path.to_str().unwrap(), + None, + HeaderMap::new(), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) .await { Ok(mut reader) => { @@ -460,10 +469,16 @@ impl DataUsageCache { match err { Error::FileNotFound | Error::VolumeNotFound => { match store - .get_object_reader(RUSTFS_META_BUCKET, name, None, HeaderMap::new(), &ObjectOptions { - no_lock: true, - ..Default::default() - }) + .get_object_reader( + RUSTFS_META_BUCKET, + name, + None, + HeaderMap::new(), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) .await { Ok(mut reader) => { @@ -804,15 +819,18 @@ impl DataUsageCache { bui.replica_count = rs.replica_count; for (arn, stat) in rs.targets.iter() { - bui.replication_info.insert(arn.clone(), BucketTargetUsageInfo { - replication_pending_size: stat.pending_size, - replicated_size: stat.replicated_size, - replication_failed_size: stat.failed_size, - replication_pending_count: stat.pending_count, - replication_failed_count: stat.failed_count, - replicated_count: stat.replicated_count, - ..Default::default() - }); + bui.replication_info.insert( + arn.clone(), + BucketTargetUsageInfo { + replication_pending_size: stat.pending_size, + replicated_size: stat.replicated_size, + replication_failed_size: stat.failed_size, + replication_pending_count: stat.pending_count, + replication_failed_count: stat.failed_count, + replicated_count: stat.replicated_count, + ..Default::default() + }, + ); } } dst.insert(bucket.name.clone(), bui); diff --git a/crates/ecstore/src/heal/heal_commands.rs b/crates/ecstore/src/heal/heal_commands.rs index 62bb6a080..3df0f3ba1 100644 --- a/crates/ecstore/src/heal/heal_commands.rs +++ b/crates/ecstore/src/heal/heal_commands.rs @@ -255,11 +255,15 @@ impl HealingTracker { pub async fn delete(&self) -> Result<()> { if let Some(disk) = &self.disk { let file_path = Path::new(BUCKET_META_PREFIX).join(HEALING_TRACKER_FILENAME); - disk.delete(RUSTFS_META_BUCKET, file_path.to_str().unwrap(), DeleteOptions { - recursive: false, - immediate: false, - ..Default::default() - }) + disk.delete( + RUSTFS_META_BUCKET, + file_path.to_str().unwrap(), + DeleteOptions { + recursive: false, + immediate: false, + ..Default::default() + }, + ) .await?; } diff --git a/crates/ecstore/src/metrics_realtime.rs b/crates/ecstore/src/metrics_realtime.rs index 0fcde38f9..298ff846f 100644 --- a/crates/ecstore/src/metrics_realtime.rs +++ b/crates/ecstore/src/metrics_realtime.rs @@ -148,11 +148,14 @@ async fn collect_local_disks_metrics(disks: &HashSet) -> HashMap res, @@ -1224,10 +1239,17 @@ impl ECStore { let mut data = PutObjReader::from_vec(chunk); let pi = match self - .put_object_part(&bucket, &object_info.name, &res.upload_id, part.number, &mut data, &ObjectOptions { - preserve_etag: Some(part.etag.clone()), - ..Default::default() - }) + .put_object_part( + &bucket, + &object_info.name, + &res.upload_id, + part.number, + &mut data, + &ObjectOptions { + preserve_etag: Some(part.etag.clone()), + ..Default::default() + }, + ) .await { Ok(pi) => pi, @@ -1247,11 +1269,17 @@ impl ECStore { if let Err(err) = self .clone() - .complete_multipart_upload(&bucket, &object_info.name, &res.upload_id, parts, &ObjectOptions { - data_movement: true, - mod_time: object_info.mod_time, - ..Default::default() - }) + .complete_multipart_upload( + &bucket, + &object_info.name, + &res.upload_id, + parts, + &ObjectOptions { + data_movement: true, + mod_time: object_info.mod_time, + ..Default::default() + }, + ) .await { error!("decommission_object: complete_multipart_upload err {:?}", &err); @@ -1267,16 +1295,21 @@ impl ECStore { let mut data = PutObjReader::new(hrd); if let Err(err) = self - .put_object(&bucket, &object_info.name, &mut data, &ObjectOptions { - src_pool_idx: pool_idx, - data_movement: true, - version_id: object_info.version_id.as_ref().map(|v| v.to_string()), - mod_time: object_info.mod_time, - user_defined: object_info.user_defined.clone(), - preserve_etag: object_info.etag.clone(), + .put_object( + &bucket, + &object_info.name, + &mut data, + &ObjectOptions { + src_pool_idx: pool_idx, + data_movement: true, + version_id: object_info.version_id.as_ref().map(|v| v.to_string()), + mod_time: object_info.mod_time, + user_defined: object_info.user_defined.clone(), + preserve_etag: object_info.etag.clone(), - ..Default::default() - }) + ..Default::default() + }, + ) .await { error!("decommission_object: put_object err {:?}", &err); @@ -1316,31 +1349,34 @@ impl SetDisks { let cb1 = cb_func.clone(); - list_path_raw(rx, ListPathRawOptions { - disks: disks.iter().cloned().map(Some).collect(), - bucket: bucket_info.name.clone(), - path: bucket_info.prefix.clone(), - recursice: true, - min_disks: listing_quorum, - agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { - let resolver = resolver.clone(); - let cb_func = cb_func.clone(); - match entries.resolve(resolver) { - Some(entry) => { - warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name); - Box::pin(async move { - cb_func(entry).await; - }) + list_path_raw( + rx, + ListPathRawOptions { + disks: disks.iter().cloned().map(Some).collect(), + bucket: bucket_info.name.clone(), + path: bucket_info.prefix.clone(), + recursice: true, + min_disks: listing_quorum, + agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + let resolver = resolver.clone(); + let cb_func = cb_func.clone(); + match entries.resolve(resolver) { + Some(entry) => { + warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name); + Box::pin(async move { + cb_func(entry).await; + }) + } + None => { + warn!("decommission_pool: list_objects_to_decommission get none"); + Box::pin(async {}) + } } - None => { - warn!("decommission_pool: list_objects_to_decommission get none"); - Box::pin(async {}) - } - } - })), - ..Default::default() - }) + })), + ..Default::default() + }, + ) .await?; Ok(()) diff --git a/crates/ecstore/src/rebalance.rs b/crates/ecstore/src/rebalance.rs index ecad7bfd1..5ab2f1499 100644 --- a/crates/ecstore/src/rebalance.rs +++ b/crates/ecstore/src/rebalance.rs @@ -774,16 +774,20 @@ impl ECStore { let mut error = None; if version.deleted { if let Err(err) = set - .delete_object(&bucket, &version.name, ObjectOptions { - versioned: true, - version_id: version_id.clone(), - mod_time: version.mod_time, - src_pool_idx: pool_index, - data_movement: true, - delete_marker: true, - skip_decommissioned: true, - ..Default::default() - }) + .delete_object( + &bucket, + &version.name, + ObjectOptions { + versioned: true, + version_id: version_id.clone(), + mod_time: version.mod_time, + src_pool_idx: pool_index, + data_movement: true, + delete_marker: true, + skip_decommissioned: true, + ..Default::default() + }, + ) .await { if is_err_object_not_found(&err) || is_err_version_not_found(&err) || is_err_data_movement_overwrite(&err) { @@ -879,12 +883,16 @@ impl ECStore { if rebalanced == fivs.versions.len() { if let Err(err) = set - .delete_object(bucket.as_str(), &encode_dir_object(&entry.name), ObjectOptions { - delete_prefix: true, - delete_prefix_object: true, + .delete_object( + bucket.as_str(), + &encode_dir_object(&entry.name), + ObjectOptions { + delete_prefix: true, + delete_prefix_object: true, - ..Default::default() - }) + ..Default::default() + }, + ) .await { error!("rebalance_entry: delete_object err {:?}", &err); @@ -903,13 +911,17 @@ impl ECStore { if object_info.is_multipart() { let res = match self - .new_multipart_upload(&bucket, &object_info.name, &ObjectOptions { - version_id: object_info.version_id.as_ref().map(|v| v.to_string()), - user_defined: object_info.user_defined.clone(), - src_pool_idx: pool_idx, - data_movement: true, - ..Default::default() - }) + .new_multipart_upload( + &bucket, + &object_info.name, + &ObjectOptions { + version_id: object_info.version_id.as_ref().map(|v| v.to_string()), + user_defined: object_info.user_defined.clone(), + src_pool_idx: pool_idx, + data_movement: true, + ..Default::default() + }, + ) .await { Ok(res) => res, @@ -943,10 +955,17 @@ impl ECStore { let mut data = PutObjReader::from_vec(chunk); let pi = match self - .put_object_part(&bucket, &object_info.name, &res.upload_id, part.number, &mut data, &ObjectOptions { - preserve_etag: Some(part.etag.clone()), - ..Default::default() - }) + .put_object_part( + &bucket, + &object_info.name, + &res.upload_id, + part.number, + &mut data, + &ObjectOptions { + preserve_etag: Some(part.etag.clone()), + ..Default::default() + }, + ) .await { Ok(pi) => pi, @@ -964,11 +983,17 @@ impl ECStore { if let Err(err) = self .clone() - .complete_multipart_upload(&bucket, &object_info.name, &res.upload_id, parts, &ObjectOptions { - data_movement: true, - mod_time: object_info.mod_time, - ..Default::default() - }) + .complete_multipart_upload( + &bucket, + &object_info.name, + &res.upload_id, + parts, + &ObjectOptions { + data_movement: true, + mod_time: object_info.mod_time, + ..Default::default() + }, + ) .await { error!("rebalance_object: complete_multipart_upload err {:?}", &err); @@ -983,16 +1008,21 @@ impl ECStore { let mut data = PutObjReader::new(hrd); if let Err(err) = self - .put_object(&bucket, &object_info.name, &mut data, &ObjectOptions { - src_pool_idx: pool_idx, - data_movement: true, - version_id: object_info.version_id.as_ref().map(|v| v.to_string()), - mod_time: object_info.mod_time, - user_defined: object_info.user_defined.clone(), - preserve_etag: object_info.etag.clone(), + .put_object( + &bucket, + &object_info.name, + &mut data, + &ObjectOptions { + src_pool_idx: pool_idx, + data_movement: true, + version_id: object_info.version_id.as_ref().map(|v| v.to_string()), + mod_time: object_info.mod_time, + user_defined: object_info.user_defined.clone(), + preserve_etag: object_info.etag.clone(), - ..Default::default() - }) + ..Default::default() + }, + ) .await { error!("rebalance_object: put_object err {:?}", &err); @@ -1137,33 +1167,36 @@ impl SetDisks { }; let cb1 = cb.clone(); - list_path_raw(rx, ListPathRawOptions { - disks: disks.iter().cloned().map(Some).collect(), - bucket: bucket.clone(), - recursice: true, - min_disks: listing_quorum, - agreed: Some(Box::new(move |entry: MetaCacheEntry| { - info!("list_objects_to_rebalance: agreed: {:?}", &entry.name); - Box::pin(cb1(entry)) - })), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { - // let cb = cb.clone(); - let resolver = resolver.clone(); - let cb = cb.clone(); + list_path_raw( + rx, + ListPathRawOptions { + disks: disks.iter().cloned().map(Some).collect(), + bucket: bucket.clone(), + recursice: true, + min_disks: listing_quorum, + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + info!("list_objects_to_rebalance: agreed: {:?}", &entry.name); + Box::pin(cb1(entry)) + })), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + // let cb = cb.clone(); + let resolver = resolver.clone(); + let cb = cb.clone(); - match entries.resolve(resolver) { - Some(entry) => { - info!("list_objects_to_rebalance: list_objects_to_decommission get {}", &entry.name); - Box::pin(async move { cb(entry).await }) + match entries.resolve(resolver) { + Some(entry) => { + info!("list_objects_to_rebalance: list_objects_to_decommission get {}", &entry.name); + Box::pin(async move { cb(entry).await }) + } + None => { + info!("list_objects_to_rebalance: list_objects_to_decommission get none"); + Box::pin(async {}) + } } - None => { - info!("list_objects_to_rebalance: list_objects_to_decommission get none"); - Box::pin(async {}) - } - } - })), - ..Default::default() - }) + })), + ..Default::default() + }, + ) .await?; info!("list_objects_to_rebalance: list_objects_to_rebalance done"); diff --git a/crates/ecstore/src/rpc/tonic_service.rs b/crates/ecstore/src/rpc/tonic_service.rs index 21ec9eb31..faaa3726b 100644 --- a/crates/ecstore/src/rpc/tonic_service.rs +++ b/crates/ecstore/src/rpc/tonic_service.rs @@ -262,10 +262,13 @@ impl Node for NodeService { let request = request.into_inner(); match self .local_peer - .delete_bucket(&request.bucket, &DeleteBucketOptions { - force: false, - ..Default::default() - }) + .delete_bucket( + &request.bucket, + &DeleteBucketOptions { + force: false, + ..Default::default() + }, + ) .await { Ok(_) => Ok(tonic::Response::new(DeleteBucketResponse { diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 91b2b539f..a0b619074 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -406,11 +406,17 @@ impl SetDisks { let src_object = src_object.clone(); futures.push(tokio::spawn(async move { let _ = disk - .delete_version(&src_bucket, &src_object, fi, false, DeleteOptions { - undo_write: true, - old_data_dir, - ..Default::default() - }) + .delete_version( + &src_bucket, + &src_object, + fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir, + ..Default::default() + }, + ) .await .map_err(|e| { debug!("rename_data delete_version err {:?}", e); @@ -483,10 +489,14 @@ impl SetDisks { tokio::spawn(async move { if let Some(disk) = disk { (disk - .delete(&bucket, &file_path, DeleteOptions { - recursive: true, - ..Default::default() - }) + .delete( + &bucket, + &file_path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) .await) .err() } else { @@ -685,10 +695,14 @@ impl SetDisks { if let Some(disk) = disks[i].as_ref() { let _ = disk - .delete(bucket, &path_join_buf(&[prefix, STORAGE_FORMAT_FILE]), DeleteOptions { - recursive: true, - ..Default::default() - }) + .delete( + bucket, + &path_join_buf(&[prefix, STORAGE_FORMAT_FILE]), + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) .await .map_err(|e| { warn!("write meta revert err {:?}", e); @@ -1667,10 +1681,14 @@ impl SetDisks { for disk in disks.iter() { futures.push(async move { if let Some(disk) = disk { - disk.delete(bucket, prefix, DeleteOptions { - recursive: true, - ..Default::default() - }) + disk.delete( + bucket, + prefix, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) .await } else { Err(DiskError::DiskNotFound) @@ -2391,10 +2409,17 @@ impl SetDisks { // Allow for dangling deletes, on versions that have DataDir missing etc. // this would end up restoring the correct readable versions. return match self - .delete_if_dang_ling(bucket, object, &parts_metadata, &errs, &data_errs_by_part, ObjectOptions { - version_id: version_id_op.clone(), - ..Default::default() - }) + .delete_if_dang_ling( + bucket, + object, + &parts_metadata, + &errs, + &data_errs_by_part, + ObjectOptions { + version_id: version_id_op.clone(), + ..Default::default() + }, + ) .await { Ok(m) => { @@ -2700,11 +2725,15 @@ impl SetDisks { if parts_metadata[index].is_remote() { let rm_data_dir = parts_metadata[index].data_dir.unwrap().to_string(); let d_path = Path::new(&encode_dir_object(object)).join(rm_data_dir); - disk.delete(bucket, d_path.to_str().unwrap(), DeleteOptions { - immediate: true, - recursive: true, - ..Default::default() - }) + disk.delete( + bucket, + d_path.to_str().unwrap(), + DeleteOptions { + immediate: true, + recursive: true, + ..Default::default() + }, + ) .await?; } @@ -2727,10 +2756,17 @@ impl SetDisks { Err(err) => { let data_errs_by_part = HashMap::new(); match self - .delete_if_dang_ling(bucket, object, &parts_metadata, &errs, &data_errs_by_part, ObjectOptions { - version_id: version_id_op.clone(), - ..Default::default() - }) + .delete_if_dang_ling( + bucket, + object, + &parts_metadata, + &errs, + &data_errs_by_part, + ObjectOptions { + version_id: version_id_op.clone(), + ..Default::default() + }, + ) .await { Ok(m) => { @@ -2786,11 +2822,15 @@ impl SetDisks { let object = object.to_string(); futures.push(tokio::spawn(async move { let _ = disk - .delete(&bucket, &object, DeleteOptions { - recursive: false, - immediate: false, - ..Default::default() - }) + .delete( + &bucket, + &object, + DeleteOptions { + recursive: false, + immediate: false, + ..Default::default() + }, + ) .await; })); } @@ -3485,11 +3525,16 @@ impl SetDisks { Ok(fivs) => fivs, Err(err) => { match self_clone - .heal_object(&bucket, &encoded_entry_name, "", &HealOpts { - scan_mode, - remove: HEAL_DELETE_DANGLING, - ..Default::default() - }) + .heal_object( + &bucket, + &encoded_entry_name, + "", + &HealOpts { + scan_mode, + remove: HEAL_DELETE_DANGLING, + ..Default::default() + }, + ) .await { Ok((res, None)) => { @@ -3606,53 +3651,56 @@ impl SetDisks { let bucket_partial = bucket.clone(); let heal_entry_agree = heal_entry.clone(); let heal_entry_partial = heal_entry.clone(); - if let Err(err) = list_path_raw(rx, ListPathRawOptions { - disks, - fallback_disks, - bucket: bucket.clone(), - recursice: true, - forward_to, - min_disks: 1, - report_not_found: false, - agreed: Some(Box::new(move |entry: MetaCacheEntry| { - let jt = jt_agree.clone(); - let bucket = bucket_agree.clone(); - let heal_entry = heal_entry_agree.clone(); - Box::pin(async move { - jt.take().await; - let bucket = bucket.clone(); - tokio::spawn(async move { - heal_entry(bucket, entry).await; - }); - }) - })), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { - let jt = jt_partial.clone(); - let bucket = bucket_partial.clone(); - let heal_entry = heal_entry_partial.clone(); - Box::pin({ - let heal_entry = heal_entry.clone(); - let resolver = resolver.clone(); - async move { - let entry = if let Some(entry) = entries.resolve(resolver) { - entry - } else if let (Some(entry), _) = entries.first_found() { - entry - } else { - return; - }; + if let Err(err) = list_path_raw( + rx, + ListPathRawOptions { + disks, + fallback_disks, + bucket: bucket.clone(), + recursice: true, + forward_to, + min_disks: 1, + report_not_found: false, + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + let jt = jt_agree.clone(); + let bucket = bucket_agree.clone(); + let heal_entry = heal_entry_agree.clone(); + Box::pin(async move { jt.take().await; let bucket = bucket.clone(); - let heal_entry = heal_entry.clone(); tokio::spawn(async move { heal_entry(bucket, entry).await; }); - } - }) - })), - finished: None, - ..Default::default() - }) + }) + })), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + let jt = jt_partial.clone(); + let bucket = bucket_partial.clone(); + let heal_entry = heal_entry_partial.clone(); + Box::pin({ + let heal_entry = heal_entry.clone(); + let resolver = resolver.clone(); + async move { + let entry = if let Some(entry) = entries.resolve(resolver) { + entry + } else if let (Some(entry), _) = entries.first_found() { + entry + } else { + return; + }; + jt.take().await; + let bucket = bucket.clone(); + let heal_entry = heal_entry.clone(); + tokio::spawn(async move { + heal_entry(bucket, entry).await; + }); + } + }) + })), + finished: None, + ..Default::default() + }, + ) .await { ret_err = Some(err.into()); @@ -3695,11 +3743,15 @@ impl SetDisks { let prefix = prefix.to_string(); futures.push(async move { if let Some(disk) = disk_op { - disk.delete(&bucket, &prefix, DeleteOptions { - recursive: true, - immediate: true, - ..Default::default() - }) + disk.delete( + &bucket, + &prefix, + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) .await } else { Ok(()) diff --git a/crates/ecstore/src/sets.rs b/crates/ecstore/src/sets.rs index da86a8309..28109afc9 100644 --- a/crates/ecstore/src/sets.rs +++ b/crates/ecstore/src/sets.rs @@ -557,11 +557,14 @@ impl StorageAPI for Sets { let idx = self.get_hashed_set_index(obj.object_name.as_str()); if !set_obj_map.contains_key(&idx) { - set_obj_map.insert(idx, vec![DelObj { - // set_idx: idx, - orig_idx: i, - obj: obj.clone(), - }]); + set_obj_map.insert( + idx, + vec![DelObj { + // set_idx: idx, + orig_idx: i, + obj: obj.clone(), + }], + ); } else if let Some(val) = set_obj_map.get_mut(&idx) { val.push(DelObj { // set_idx: idx, @@ -753,10 +756,13 @@ impl StorageAPI for Sets { #[tracing::instrument(skip(self))] async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { - let (disks, _) = init_storage_disks_with_errors(&self.endpoints.endpoints, &DiskOption { - cleanup: false, - health_check: false, - }) + let (disks, _) = init_storage_disks_with_errors( + &self.endpoints.endpoints, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) .await; let (formats, errs) = load_format_erasure_all(&disks, true).await; if let Err(err) = check_format_erasure_values(&formats, self.set_drive_count) { diff --git a/crates/ecstore/src/store.rs b/crates/ecstore/src/store.rs index 5e69333d7..bff22f7d2 100644 --- a/crates/ecstore/src/store.rs +++ b/crates/ecstore/src/store.rs @@ -154,10 +154,13 @@ impl ECStore { // validate_parity(partiy_count, pool_eps.drives_per_set)?; - let (disks, errs) = store_init::init_disks(&pool_eps.endpoints, &DiskOption { - cleanup: true, - health_check: true, - }) + let (disks, errs) = store_init::init_disks( + &pool_eps.endpoints, + &DiskOption { + cleanup: true, + health_check: true, + }, + ) .await; check_disk_fatal_errs(&errs)?; @@ -501,10 +504,14 @@ impl ECStore { } async fn delete_prefix(&self, bucket: &str, object: &str) -> Result<()> { for pool in self.pools.iter() { - pool.delete_object(bucket, object, ObjectOptions { - delete_prefix: true, - ..Default::default() - }) + pool.delete_object( + bucket, + object, + ObjectOptions { + delete_prefix: true, + ..Default::default() + }, + ) .await?; } @@ -614,11 +621,15 @@ impl ECStore { async fn get_pool_idx(&self, bucket: &str, object: &str, size: i64) -> Result { let idx = match self - .get_pool_idx_existing_with_opts(bucket, object, &ObjectOptions { - skip_decommissioned: true, - skip_rebalancing: true, - ..Default::default() - }) + .get_pool_idx_existing_with_opts( + bucket, + object, + &ObjectOptions { + skip_decommissioned: true, + skip_rebalancing: true, + ..Default::default() + }, + ) .await { Ok(res) => res, @@ -659,12 +670,16 @@ impl ECStore { } async fn get_pool_idx_existing_no_lock(&self, bucket: &str, object: &str) -> Result { - self.get_pool_idx_existing_with_opts(bucket, object, &ObjectOptions { - no_lock: true, - skip_decommissioned: true, - skip_rebalancing: true, - ..Default::default() - }) + self.get_pool_idx_existing_with_opts( + bucket, + object, + &ObjectOptions { + no_lock: true, + skip_decommissioned: true, + skip_rebalancing: true, + ..Default::default() + }, + ) .await } @@ -1359,11 +1374,14 @@ impl StorageAPI for ECStore { if let Err(err) = self.peer_sys.make_bucket(bucket, opts).await { if !is_err_bucket_exists(&err.into()) { let _ = self - .delete_bucket(bucket, &DeleteBucketOptions { - no_lock: true, - no_recreate: true, - ..Default::default() - }) + .delete_bucket( + bucket, + &DeleteBucketOptions { + no_lock: true, + no_recreate: true, + ..Default::default() + }, + ) .await; } }; @@ -1666,10 +1684,14 @@ impl StorageAPI for ECStore { for obj in objects.iter() { futures.push(async move { - self.internal_get_pool_info_existing_with_opts(bucket, &obj.object_name, &ObjectOptions { - no_lock: true, - ..Default::default() - }) + self.internal_get_pool_info_existing_with_opts( + bucket, + &obj.object_name, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) .await }); } diff --git a/crates/ecstore/src/store_list_objects.rs b/crates/ecstore/src/store_list_objects.rs index 8c9e0c3dc..fa6c4d7f7 100644 --- a/crates/ecstore/src/store_list_objects.rs +++ b/crates/ecstore/src/store_list_objects.rs @@ -263,10 +263,14 @@ impl ECStore { // use get if !opts.prefix.is_empty() && opts.limit == 1 && opts.marker.is_none() { match self - .get_object_info(&opts.bucket, &opts.prefix, &ObjectOptions { - no_lock: true, - ..Default::default() - }) + .get_object_info( + &opts.bucket, + &opts.prefix, + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) .await { Ok(res) => { @@ -765,45 +769,48 @@ impl ECStore { let tx1 = sender.clone(); let tx2 = sender.clone(); - list_path_raw(rx.resubscribe(), ListPathRawOptions { - disks: disks.iter().cloned().map(Some).collect(), - fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), - bucket: bucket.to_owned(), - path, - recursice: true, - filter_prefix: Some(filter_prefix), - forward_to: opts.marker.clone(), - min_disks: listing_quorum, - per_disk_limit: opts.limit as i32, - agreed: Some(Box::new(move |entry: MetaCacheEntry| { - Box::pin({ - let value = tx1.clone(); - async move { - if entry.is_dir() { - return; - } - if let Err(err) = value.send(entry).await { - error!("list_path send fail {:?}", err); - } - } - }) - })), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { - Box::pin({ - let value = tx2.clone(); - let resolver = resolver.clone(); - async move { - if let Some(entry) = entries.resolve(resolver) { + list_path_raw( + rx.resubscribe(), + ListPathRawOptions { + disks: disks.iter().cloned().map(Some).collect(), + fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), + bucket: bucket.to_owned(), + path, + recursice: true, + filter_prefix: Some(filter_prefix), + forward_to: opts.marker.clone(), + min_disks: listing_quorum, + per_disk_limit: opts.limit as i32, + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + Box::pin({ + let value = tx1.clone(); + async move { + if entry.is_dir() { + return; + } if let Err(err) = value.send(entry).await { error!("list_path send fail {:?}", err); } } - } - }) - })), - finished: None, - ..Default::default() - }) + }) + })), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + Box::pin({ + let value = tx2.clone(); + let resolver = resolver.clone(); + async move { + if let Some(entry) = entries.resolve(resolver) { + if let Err(err) = value.send(entry).await { + error!("list_path send fail {:?}", err); + } + } + } + }) + })), + finished: None, + ..Default::default() + }, + ) .await }); } @@ -1272,42 +1279,45 @@ impl SetDisks { let tx1 = sender.clone(); let tx2 = sender.clone(); - list_path_raw(rx, ListPathRawOptions { - disks: disks.iter().cloned().map(Some).collect(), - fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), - bucket: opts.bucket, - path: opts.base_dir, - recursice: opts.recursive, - filter_prefix: opts.filter_prefix, - forward_to: opts.marker, - min_disks: listing_quorum, - per_disk_limit: limit, - agreed: Some(Box::new(move |entry: MetaCacheEntry| { - Box::pin({ - let value = tx1.clone(); - async move { - if let Err(err) = value.send(entry).await { - error!("list_path send fail {:?}", err); - } - } - }) - })), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { - Box::pin({ - let value = tx2.clone(); - let resolver = resolver.clone(); - async move { - if let Some(entry) = entries.resolve(resolver) { + list_path_raw( + rx, + ListPathRawOptions { + disks: disks.iter().cloned().map(Some).collect(), + fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), + bucket: opts.bucket, + path: opts.base_dir, + recursice: opts.recursive, + filter_prefix: opts.filter_prefix, + forward_to: opts.marker, + min_disks: listing_quorum, + per_disk_limit: limit, + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + Box::pin({ + let value = tx1.clone(); + async move { if let Err(err) = value.send(entry).await { error!("list_path send fail {:?}", err); } } - } - }) - })), - finished: None, - ..Default::default() - }) + }) + })), + partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { + Box::pin({ + let value = tx2.clone(); + let resolver = resolver.clone(); + async move { + if let Some(entry) = entries.resolve(resolver) { + if let Err(err) = value.send(entry).await { + error!("list_path send fail {:?}", err); + } + } + } + }) + })), + finished: None, + ..Default::default() + }, + ) .await .map_err(Error::other) } diff --git a/crates/ecstore/src/tier/tier.rs b/crates/ecstore/src/tier/tier.rs index 9293ff152..fdfed4983 100644 --- a/crates/ecstore/src/tier/tier.rs +++ b/crates/ecstore/src/tier/tier.rs @@ -365,10 +365,15 @@ impl TierConfigMgr { file: &str, data: Bytes, ) -> std::result::Result<(), std::io::Error> { - self.save_config_with_opts(api, file, data, &ObjectOptions { - max_parity: true, - ..Default::default() - }) + self.save_config_with_opts( + api, + file, + data, + &ObjectOptions { + max_parity: true, + ..Default::default() + }, + ) .await } diff --git a/crates/ecstore/src/tier/warm_backend_minio.rs b/crates/ecstore/src/tier/warm_backend_minio.rs index b8eeccfe7..73da4acf4 100644 --- a/crates/ecstore/src/tier/warm_backend_minio.rs +++ b/crates/ecstore/src/tier/warm_backend_minio.rs @@ -101,13 +101,19 @@ impl WarmBackend for WarmBackendMinIO { let part_size = optimal_part_size(length)?; let client = self.0.client.clone(); let res = client - .put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &PutObjectOptions { - storage_class: self.0.storage_class.clone(), - part_size: part_size as u64, - disable_content_sha256: true, - user_metadata: meta, - ..Default::default() - }) + .put_object( + &self.0.bucket, + &self.0.get_dest(object), + r, + length, + &PutObjectOptions { + storage_class: self.0.storage_class.clone(), + part_size: part_size as u64, + disable_content_sha256: true, + user_metadata: meta, + ..Default::default() + }, + ) .await?; //self.ToObjectError(err, object) Ok(res.version_id) diff --git a/crates/ecstore/src/tier/warm_backend_rustfs.rs b/crates/ecstore/src/tier/warm_backend_rustfs.rs index 79e79a90d..8bc8142b3 100644 --- a/crates/ecstore/src/tier/warm_backend_rustfs.rs +++ b/crates/ecstore/src/tier/warm_backend_rustfs.rs @@ -98,13 +98,19 @@ impl WarmBackend for WarmBackendRustFS { let part_size = optimal_part_size(length)?; let client = self.0.client.clone(); let res = client - .put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &PutObjectOptions { - storage_class: self.0.storage_class.clone(), - part_size: part_size as u64, - disable_content_sha256: true, - user_metadata: meta, - ..Default::default() - }) + .put_object( + &self.0.bucket, + &self.0.get_dest(object), + r, + length, + &PutObjectOptions { + storage_class: self.0.storage_class.clone(), + part_size: part_size as u64, + disable_content_sha256: true, + user_metadata: meta, + ..Default::default() + }, + ) .await?; //self.ToObjectError(err, object) Ok(res.version_id) diff --git a/crates/ecstore/src/tier/warm_backend_s3.rs b/crates/ecstore/src/tier/warm_backend_s3.rs index 7609e6fd8..f6de4b75d 100644 --- a/crates/ecstore/src/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/tier/warm_backend_s3.rs @@ -127,12 +127,18 @@ impl WarmBackend for WarmBackendS3 { ) -> Result { let client = self.client.clone(); let res = client - .put_object(&self.bucket, &self.get_dest(object), r, length, &PutObjectOptions { - send_content_md5: true, - storage_class: self.storage_class.clone(), - user_metadata: meta, - ..Default::default() - }) + .put_object( + &self.bucket, + &self.get_dest(object), + r, + length, + &PutObjectOptions { + send_content_md5: true, + storage_class: self.storage_class.clone(), + user_metadata: meta, + ..Default::default() + }, + ) .await?; Ok(res.version_id) } diff --git a/crates/iam/src/manager.rs b/crates/iam/src/manager.rs index a2bdf0b95..c0831ce11 100644 --- a/crates/iam/src/manager.rs +++ b/crates/iam/src/manager.rs @@ -1575,11 +1575,14 @@ pub fn get_default_policyes() -> HashMap { default_policies .iter() .map(|(n, p)| { - (n.to_string(), PolicyDoc { - version: 1, - policy: p.clone(), - ..Default::default() - }) + ( + n.to_string(), + PolicyDoc { + version: 1, + policy: p.clone(), + ..Default::default() + }, + ) }) .collect() } diff --git a/crates/lock/src/local_locker.rs b/crates/lock/src/local_locker.rs index 8a65d5b64..ef9676c5c 100644 --- a/crates/lock/src/local_locker.rs +++ b/crates/lock/src/local_locker.rs @@ -140,17 +140,20 @@ impl Locker for LocalLocker { } args.resources.iter().enumerate().for_each(|(idx, resource)| { - self.lock_map.insert(resource.to_string(), vec![LockRequesterInfo { - name: resource.to_string(), - writer: true, - source: args.source.to_string(), - owner: args.owner.to_string(), - uid: args.uid.to_string(), - group: args.resources.len() > 1, - quorum: args.quorum, - idx, - ..Default::default() - }]); + self.lock_map.insert( + resource.to_string(), + vec![LockRequesterInfo { + name: resource.to_string(), + writer: true, + source: args.source.to_string(), + owner: args.owner.to_string(), + uid: args.uid.to_string(), + group: args.resources.len() > 1, + quorum: args.quorum, + idx, + ..Default::default() + }], + ); let mut uuid = args.uid.to_string(); format_uuid(&mut uuid, &idx); @@ -227,15 +230,18 @@ impl Locker for LocalLocker { } } None => { - self.lock_map.insert(resource.to_string(), vec![LockRequesterInfo { - name: resource.to_string(), - writer: false, - source: args.source.to_string(), - owner: args.owner.to_string(), - uid: args.uid.to_string(), - quorum: args.quorum, - ..Default::default() - }]); + self.lock_map.insert( + resource.to_string(), + vec![LockRequesterInfo { + name: resource.to_string(), + writer: false, + source: args.source.to_string(), + owner: args.owner.to_string(), + uid: args.uid.to_string(), + quorum: args.quorum, + ..Default::default() + }], + ); } } let mut uuid = args.uid.to_string(); diff --git a/crates/obs/examples/server.rs b/crates/obs/examples/server.rs index f1bea80c6..fc413957a 100644 --- a/crates/obs/examples/server.rs +++ b/crates/obs/examples/server.rs @@ -46,10 +46,10 @@ async fn run(bucket: String, object: String, user: String, service_name: String) // Record Metrics let meter = global::meter("rustfs"); let request_duration = meter.f64_histogram("s3_request_duration_seconds").build(); - request_duration.record(start_time.elapsed().unwrap().as_secs_f64(), &[opentelemetry::KeyValue::new( - "operation", - "run", - )]); + request_duration.record( + start_time.elapsed().unwrap().as_secs_f64(), + &[opentelemetry::KeyValue::new("operation", "run")], + ); match SystemObserver::init_process_observer(meter).await { Ok(_) => info!("Process observer initialized successfully"), @@ -84,10 +84,10 @@ async fn put_object(bucket: String, object: String, user: String) { let meter = global::meter("rustfs"); let request_duration = meter.f64_histogram("s3_request_duration_seconds").build(); - request_duration.record(start_time.elapsed().unwrap().as_secs_f64(), &[opentelemetry::KeyValue::new( - "operation", - "put_object", - )]); + request_duration.record( + start_time.elapsed().unwrap().as_secs_f64(), + &[opentelemetry::KeyValue::new("operation", "put_object")], + ); info!( "Starting PUT operation content: bucket = {}, object = {}, user = {},start_time = {}", diff --git a/crates/obs/src/system/collector.rs b/crates/obs/src/system/collector.rs index 68f628dd0..ecaa05101 100644 --- a/crates/obs/src/system/collector.rs +++ b/crates/obs/src/system/collector.rs @@ -115,18 +115,24 @@ impl Collector { let transmitted = data.transmitted() as i64; self.metrics.network_io_per_interface.record( received, - &[&self.attributes.attributes[..], &[ - KeyValue::new(INTERFACE, interface_name.to_string()), - KeyValue::new(DIRECTION, "received"), - ]] + &[ + &self.attributes.attributes[..], + &[ + KeyValue::new(INTERFACE, interface_name.to_string()), + KeyValue::new(DIRECTION, "received"), + ], + ] .concat(), ); self.metrics.network_io_per_interface.record( transmitted, - &[&self.attributes.attributes[..], &[ - KeyValue::new(INTERFACE, interface_name.to_string()), - KeyValue::new(DIRECTION, "transmitted"), - ]] + &[ + &self.attributes.attributes[..], + &[ + KeyValue::new(INTERFACE, interface_name.to_string()), + KeyValue::new(DIRECTION, "transmitted"), + ], + ] .concat(), ); } @@ -149,10 +155,10 @@ impl Collector { }; self.metrics.process_status.record( status_value, - &[&self.attributes.attributes[..], &[KeyValue::new( - STATUS, - format!("{:?}", process.status()), - )]] + &[ + &self.attributes.attributes[..], + &[KeyValue::new(STATUS, format!("{:?}", process.status()))], + ] .concat(), ); diff --git a/crates/policy/src/policy/policy.rs b/crates/policy/src/policy/policy.rs index 01ef5181d..1a3f5ef2f 100644 --- a/crates/policy/src/policy/policy.rs +++ b/crates/policy/src/policy/policy.rs @@ -276,150 +276,12 @@ pub mod default { #[allow(clippy::incompatible_msrv)] pub static DEFAULT_POLICIES: LazyLock<[(&'static str, Policy); 6]> = LazyLock::new(|| { [ - ("readwrite", Policy { - id: "".into(), - version: DEFAULT_VERSION.into(), - statements: vec![Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::S3Action(S3Action::AllActions)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Resource::S3("*".into())); - hash_set - }), - conditions: Functions::default(), - ..Default::default() - }], - }), - ("readonly", Policy { - id: "".into(), - version: DEFAULT_VERSION.into(), - statements: vec![Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::S3Action(S3Action::GetBucketLocationAction)); - hash_set.insert(Action::S3Action(S3Action::GetObjectAction)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Resource::S3("*".into())); - hash_set - }), - conditions: Functions::default(), - ..Default::default() - }], - }), - ("writeonly", Policy { - id: "".into(), - version: DEFAULT_VERSION.into(), - statements: vec![Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::S3Action(S3Action::PutObjectAction)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Resource::S3("*".into())); - hash_set - }), - conditions: Functions::default(), - ..Default::default() - }], - }), - ("writeonly", Policy { - id: "".into(), - version: DEFAULT_VERSION.into(), - statements: vec![Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::S3Action(S3Action::PutObjectAction)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Resource::S3("*".into())); - hash_set - }), - conditions: Functions::default(), - ..Default::default() - }], - }), - ("diagnostics", Policy { - id: "".into(), - version: DEFAULT_VERSION.into(), - statements: vec![Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::AdminAction(AdminAction::ProfilingAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::TraceAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::ConsoleLogAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::ServerInfoAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::TopLocksAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::HealthInfoAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::PrometheusAdminAction)); - hash_set.insert(Action::AdminAction(AdminAction::BandwidthMonitorAction)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Resource::S3("*".into())); - hash_set - }), - conditions: Functions::default(), - ..Default::default() - }], - }), - ("consoleAdmin", Policy { - id: "".into(), - version: DEFAULT_VERSION.into(), - statements: vec![ - Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::AdminAction(AdminAction::AllAdminActions)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet(HashSet::new()), - conditions: Functions::default(), - ..Default::default() - }, - Statement { - sid: "".into(), - effect: Effect::Allow, - actions: ActionSet({ - let mut hash_set = HashSet::new(); - hash_set.insert(Action::KmsAction(KmsAction::AllActions)); - hash_set - }), - not_actions: ActionSet(Default::default()), - resources: ResourceSet(HashSet::new()), - conditions: Functions::default(), - ..Default::default() - }, - Statement { + ( + "readwrite", + Policy { + id: "".into(), + version: DEFAULT_VERSION.into(), + statements: vec![Statement { sid: "".into(), effect: Effect::Allow, actions: ActionSet({ @@ -435,9 +297,165 @@ pub mod default { }), conditions: Functions::default(), ..Default::default() - }, - ], - }), + }], + }, + ), + ( + "readonly", + Policy { + id: "".into(), + version: DEFAULT_VERSION.into(), + statements: vec![Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::S3Action(S3Action::GetBucketLocationAction)); + hash_set.insert(Action::S3Action(S3Action::GetObjectAction)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Resource::S3("*".into())); + hash_set + }), + conditions: Functions::default(), + ..Default::default() + }], + }, + ), + ( + "writeonly", + Policy { + id: "".into(), + version: DEFAULT_VERSION.into(), + statements: vec![Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::S3Action(S3Action::PutObjectAction)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Resource::S3("*".into())); + hash_set + }), + conditions: Functions::default(), + ..Default::default() + }], + }, + ), + ( + "writeonly", + Policy { + id: "".into(), + version: DEFAULT_VERSION.into(), + statements: vec![Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::S3Action(S3Action::PutObjectAction)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Resource::S3("*".into())); + hash_set + }), + conditions: Functions::default(), + ..Default::default() + }], + }, + ), + ( + "diagnostics", + Policy { + id: "".into(), + version: DEFAULT_VERSION.into(), + statements: vec![Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::AdminAction(AdminAction::ProfilingAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::TraceAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::ConsoleLogAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::ServerInfoAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::TopLocksAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::HealthInfoAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::PrometheusAdminAction)); + hash_set.insert(Action::AdminAction(AdminAction::BandwidthMonitorAction)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Resource::S3("*".into())); + hash_set + }), + conditions: Functions::default(), + ..Default::default() + }], + }, + ), + ( + "consoleAdmin", + Policy { + id: "".into(), + version: DEFAULT_VERSION.into(), + statements: vec![ + Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::AdminAction(AdminAction::AllAdminActions)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet(HashSet::new()), + conditions: Functions::default(), + ..Default::default() + }, + Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::KmsAction(KmsAction::AllActions)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet(HashSet::new()), + conditions: Functions::default(), + ..Default::default() + }, + Statement { + sid: "".into(), + effect: Effect::Allow, + actions: ActionSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Action::S3Action(S3Action::AllActions)); + hash_set + }), + not_actions: ActionSet(Default::default()), + resources: ResourceSet({ + let mut hash_set = HashSet::new(); + hash_set.insert(Resource::S3("*".into())); + hash_set + }), + conditions: Functions::default(), + ..Default::default() + }, + ], + }, + ), ] }); } diff --git a/crates/protos/src/main.rs b/crates/protos/src/main.rs index 0a7118b2b..5e48c4d57 100644 --- a/crates/protos/src/main.rs +++ b/crates/protos/src/main.rs @@ -85,9 +85,13 @@ fn main() -> Result<(), AnyError> { Err(_) => "flatc".to_string(), }; - compile_flatbuffers_models(&mut generated_mod_rs, &flatc_path, proto_dir.clone(), flatbuffer_out_dir.clone(), vec![ - "models", - ])?; + compile_flatbuffers_models( + &mut generated_mod_rs, + &flatc_path, + proto_dir.clone(), + flatbuffer_out_dir.clone(), + vec!["models"], + )?; fmt(); Ok(()) diff --git a/crates/signer/src/request_signature_v4.rs b/crates/signer/src/request_signature_v4.rs index 5f89df8ce..7b2f14c0e 100644 --- a/crates/signer/src/request_signature_v4.rs +++ b/crates/signer/src/request_signature_v4.rs @@ -340,10 +340,7 @@ fn sign_v4_inner( let headers = req.headers_mut(); - let auth = format!( - "{} Credential={}, SignedHeaders={}, Signature={}", - SIGN_V4_ALGORITHM, credential, signed_headers, signature - ); + let auth = format!("{SIGN_V4_ALGORITHM} Credential={credential}, SignedHeaders={signed_headers}, Signature={signature}"); headers.insert("Authorization", auth.parse().unwrap()); if !trailer.is_empty() { diff --git a/rust-toolchain.toml b/rust-toolchain.toml new file mode 100644 index 000000000..86cba4f71 --- /dev/null +++ b/rust-toolchain.toml @@ -0,0 +1,17 @@ +# Copyright 2024 RustFS Team +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +[toolchain] +channel = "stable" +components = ["rustfmt", "clippy", "rust-src", "rust-analyzer"] diff --git a/rustfs/src/admin/handlers/bucket_meta.rs b/rustfs/src/admin/handlers/bucket_meta.rs index a5f0fea8b..074349bc9 100644 --- a/rustfs/src/admin/handlers/bucket_meta.rs +++ b/rustfs/src/admin/handlers/bucket_meta.rs @@ -453,10 +453,13 @@ impl Operation for ImportBucketMetadata { // create bucket if not exists if !bucket_metadatas.contains_key(bucket_name) { if let Err(e) = store - .make_bucket(bucket_name, &MakeBucketOptions { - force_create: true, - ..Default::default() - }) + .make_bucket( + bucket_name, + &MakeBucketOptions { + force_create: true, + ..Default::default() + }, + ) .await { warn!("create bucket failed: {e}"); diff --git a/rustfs/src/admin/handlers/user.rs b/rustfs/src/admin/handlers/user.rs index 60c978594..64ebcb41a 100644 --- a/rustfs/src/admin/handlers/user.rs +++ b/rustfs/src/admin/handlers/user.rs @@ -445,17 +445,20 @@ impl Operation for ExportIam { let users: HashMap = users .into_iter() .map(|(k, v)| { - (k, AddOrUpdateUserReq { - secret_key: v.credentials.secret_key, - status: { - if v.credentials.status == "off" { - AccountStatus::Disabled - } else { - AccountStatus::Enabled - } + ( + k, + AddOrUpdateUserReq { + secret_key: v.credentials.secret_key, + status: { + if v.credentials.status == "off" { + AccountStatus::Disabled + } else { + AccountStatus::Enabled + } + }, + policy: None, }, - policy: None, - }) + ) }) .collect::>(); diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index a6b524875..6782f5c09 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -310,11 +310,14 @@ impl S3 for FS { }; store - .make_bucket(&bucket, &MakeBucketOptions { - force_create: true, - lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v), - ..Default::default() - }) + .make_bucket( + &bucket, + &MakeBucketOptions { + force_create: true, + lock_enabled: object_lock_enabled_for_bucket.is_some_and(|v| v), + ..Default::default() + }, + ) .await .map_err(ApiError::from)?; @@ -495,10 +498,13 @@ impl S3 for FS { }; store - .delete_bucket(&input.bucket, &DeleteBucketOptions { - force: false, - ..Default::default() - }) + .delete_bucket( + &input.bucket, + &DeleteBucketOptions { + force: false, + ..Default::default() + }, + ) .await .map_err(ApiError::from)?;