fix: some expect error (#2622)

Signed-off-by: likewu <likewu@126.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
likewu
2026-04-21 11:27:50 +08:00
committed by GitHub
parent 1525143a04
commit a77be8f89b
23 changed files with 613 additions and 395 deletions
@@ -70,7 +70,7 @@ use time::OffsetDateTime;
use tokio::select; use tokio::select;
use tokio::sync::mpsc::{Receiver, Sender}; use tokio::sync::mpsc::{Receiver, Sender};
use tokio::sync::{RwLock, mpsc}; use tokio::sync::{RwLock, mpsc};
use tracing::{error, info, warn}; use tracing::{debug, error, info, warn};
use uuid::Uuid; use uuid::Uuid;
use xxhash_rust::xxh64; use xxhash_rust::xxh64;
@@ -413,6 +413,7 @@ impl ExpiryState {
let v = v.expect("received None after None check"); let v = v.expect("received None after None check");
if v.as_any().is::<ExpiryTask>() { if v.as_any().is::<ExpiryTask>() {
let v = v.as_any().downcast_ref::<ExpiryTask>().expect("ExpiryTask downcast failed"); let v = v.as_any().downcast_ref::<ExpiryTask>().expect("ExpiryTask downcast failed");
//debug!("lifecycle expiry worker received task: {:?}", v.obj_info);
if !v.obj_info.transitioned_object.status.is_empty() { if !v.obj_info.transitioned_object.status.is_empty() {
apply_expiry_on_transitioned_object(api.clone(), &v.obj_info, &v.event, &v.src).await; apply_expiry_on_transitioned_object(api.clone(), &v.obj_info, &v.event, &v.src).await;
} else { } else {
@@ -1340,8 +1341,8 @@ pub async fn expire_transitioned_object(
&oi.transitioned_object.tier, &oi.transitioned_object.tier,
) )
.await; .await;
if ret.is_err() { if let Err(e) = &ret {
//transitionLogIf(ctx, err); error!("Failed to delete remote transitioned object {}: {:?}", oi.transitioned_object.name, e);
} }
mark_delete_opts_skip_decommissioned_on_remote_success(&mut opts, ret.is_ok()); mark_delete_opts_skip_decommissioned_on_remote_success(&mut opts, ret.is_ok());
@@ -1356,7 +1357,7 @@ pub async fn expire_transitioned_object(
schedule_lifecycle_replication_delete_if_needed(oi, &dobj).await; schedule_lifecycle_replication_delete_if_needed(oi, &dobj).await;
//defer auditLogLifecycle(ctx, *oi, ILMExpiry, tags, traceFn) //audit_log_lifecycle(oi, ILMExpiry, tags);
let event_name = if oi.delete_marker { let event_name = if oi.delete_marker {
EventName::LifecycleExpirationDelete EventName::LifecycleExpirationDelete
@@ -1770,7 +1771,7 @@ pub async fn apply_expiry_on_non_transitioned_objects(
let time_ilm = Metrics::time_ilm(lc_event.action); let time_ilm = Metrics::time_ilm(lc_event.action);
//debug!("lc_event.action: {:?}", lc_event.action); //debug!("lc_event.action: {:?}", lc_event.action);
//debug!("opts: {:?}", opts); debug!("expiry_on_non_transitioned_objects opts: {:?}", opts);
let mut dobj = match api.delete_object(&oi.bucket, &encode_dir_object(&oi.name), opts).await { let mut dobj = match api.delete_object(&oi.bucket, &encode_dir_object(&oi.name), opts).await {
Ok(dobj) => dobj, Ok(dobj) => dobj,
Err(e) => { Err(e) => {
+2 -2
View File
@@ -439,8 +439,8 @@ impl Lifecycle for BucketLifecycleConfiguration {
async fn eval_inner(&self, obj: &ObjectOpts, now: OffsetDateTime, _newer_noncurrent_versions: usize) -> Event { async fn eval_inner(&self, obj: &ObjectOpts, now: OffsetDateTime, _newer_noncurrent_versions: usize) -> Event {
let mut events = Vec::<Event>::new(); let mut events = Vec::<Event>::new();
debug!( debug!(
"eval_inner: object={}, mod_time={:?}, now={:?}, is_latest={}, delete_marker={}", "eval_inner: object={}, mod_time={:?}, successor_mod_time={:?}, now={:?}, is_latest={}, delete_marker={}",
obj.name, obj.mod_time, now, obj.is_latest, obj.delete_marker obj.name, obj.mod_time, obj.successor_mod_time, now, obj.is_latest, obj.delete_marker
); );
// Gracefully handle missing mod_time instead of panicking // Gracefully handle missing mod_time instead of panicking
@@ -1,4 +1,3 @@
#![allow(unused_imports)]
// Copyright 2024 RustFS Team // Copyright 2024 RustFS Team
// //
// Licensed under the Apache License, Version 2.0 (the "License"); // Licensed under the Apache License, Version 2.0 (the "License");
@@ -12,6 +11,7 @@
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
#![allow(unused_imports)]
#![allow(unused_variables)] #![allow(unused_variables)]
#![allow(unused_mut)] #![allow(unused_mut)]
#![allow(unused_assignments)] #![allow(unused_assignments)]
+23 -13
View File
@@ -100,7 +100,7 @@ pub fn http_resp_to_error_response(
bucket_name: &str, bucket_name: &str,
object_name: &str, object_name: &str,
) -> ErrorResponse { ) -> ErrorResponse {
let err_body = String::from_utf8(b).unwrap(); let err_body = String::from_utf8_lossy(&b).to_string();
if h.is_empty() || resp_status.is_client_error() || resp_status.is_server_error() { if h.is_empty() || resp_status.is_client_error() || resp_status.is_server_error() {
return ErrorResponse { return ErrorResponse {
status_code: resp_status, status_code: resp_status,
@@ -178,36 +178,46 @@ pub fn http_resp_to_error_response(
}; };
} }
} }
} else { } else if let Ok(parsed_resp) = err_resp_ {
err_resp = err_resp_.unwrap(); err_resp = parsed_resp;
} }
err_resp.status_code = resp_status; err_resp.status_code = resp_status;
if let Some(server_name) = h.get("Server") { if let Some(server_name) = h.get("Server") {
err_resp.server = server_name.to_str().expect("err").to_string(); if let Ok(server_str) = server_name.to_str() {
err_resp.server = server_str.to_string();
}
} }
let code = h.get("x-minio-error-code"); if let Some(code) = h.get("x-minio-error-code") {
if code.is_some() { if let Ok(code_str) = code.to_str() {
err_resp.code = S3ErrorCode::Custom(code.expect("err").to_str().expect("err").into()); err_resp.code = S3ErrorCode::Custom(code_str.into());
}
} }
let desc = h.get("x-minio-error-desc"); if let Some(desc) = h.get("x-minio-error-desc") {
if desc.is_some() { if let Ok(desc_str) = desc.to_str() {
err_resp.message = desc.expect("err").to_str().expect("err").trim_matches('"').to_string(); err_resp.message = desc_str.trim_matches('"').to_string();
}
} }
if err_resp.request_id == "" { if err_resp.request_id == "" {
if let Some(x_amz_request_id) = h.get("x-amz-request-id") { if let Some(x_amz_request_id) = h.get("x-amz-request-id") {
err_resp.request_id = x_amz_request_id.to_str().expect("err").to_string(); if let Ok(request_id_str) = x_amz_request_id.to_str() {
err_resp.request_id = request_id_str.to_string();
}
} }
} }
if err_resp.host_id == "" { if err_resp.host_id == "" {
if let Some(x_amz_id_2) = h.get("x-amz-id-2") { if let Some(x_amz_id_2) = h.get("x-amz-id-2") {
err_resp.host_id = x_amz_id_2.to_str().expect("err").to_string(); if let Ok(host_id_str) = x_amz_id_2.to_str() {
err_resp.host_id = host_id_str.to_string();
}
} }
} }
if err_resp.region == "" { if err_resp.region == "" {
if let Some(x_amz_bucket_region) = h.get("x-amz-bucket-region") { if let Some(x_amz_bucket_region) = h.get("x-amz-bucket-region") {
err_resp.region = x_amz_bucket_region.to_str().expect("err").to_string(); if let Ok(region_str) = x_amz_bucket_region.to_str() {
err_resp.region = region_str.to_string();
}
} }
} }
if err_resp.code == S3ErrorCode::InvalidLocationConstraint/*InvalidRegion*/ && err_resp.region != "" { if err_resp.code == S3ErrorCode::InvalidLocationConstraint/*InvalidRegion*/ && err_resp.region != "" {
+1 -1
View File
@@ -1,4 +1,3 @@
#![allow(clippy::map_entry)]
// Copyright 2024 RustFS Team // Copyright 2024 RustFS Team
// //
// Licensed under the Apache License, Version 2.0 (the "License"); // Licensed under the Apache License, Version 2.0 (the "License");
@@ -12,6 +11,7 @@
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
#![allow(clippy::map_entry)]
#![allow(unused_imports)] #![allow(unused_imports)]
#![allow(unused_variables)] #![allow(unused_variables)]
#![allow(unused_mut)] #![allow(unused_mut)]
+57 -37
View File
@@ -137,21 +137,21 @@ impl Default for PutObjectOptions {
impl PutObjectOptions { impl PutObjectOptions {
fn set_match_etag(&mut self, etag: &str) { fn set_match_etag(&mut self, etag: &str) {
if etag == "*" { if etag == "*" {
self.custom_header self.custom_header.insert("If-Match", HeaderValue::from_static("*"));
.insert("If-Match", HeaderValue::from_str("*").expect("err"));
} else { } else {
self.custom_header if let Ok(etag_value) = HeaderValue::from_str(&format!("\"{}\"", etag)) {
.insert("If-Match", HeaderValue::from_str(&format!("\"{}\"", etag)).expect("err")); self.custom_header.insert("If-Match", etag_value);
}
} }
} }
fn set_match_etag_except(&mut self, etag: &str) { fn set_match_etag_except(&mut self, etag: &str) {
if etag == "*" { if etag == "*" {
self.custom_header self.custom_header.insert("If-None-Match", HeaderValue::from_static("*"));
.insert("If-None-Match", HeaderValue::from_str("*").expect("err"));
} else { } else {
self.custom_header if let Ok(etag_value) = HeaderValue::from_str(&format!("\"{etag}\"")) {
.insert("If-None-Match", HeaderValue::from_str(&format!("\"{etag}\"")).expect("err")); self.custom_header.insert("If-None-Match", etag_value);
}
} }
} }
@@ -162,59 +162,75 @@ impl PutObjectOptions {
if content_type == "" { if content_type == "" {
content_type = "application/octet-stream".to_string(); content_type = "application/octet-stream".to_string();
} }
header.insert("Content-Type", HeaderValue::from_str(&content_type).expect("err")); if let Ok(content_type_value) = HeaderValue::from_str(&content_type) {
header.insert("Content-Type", content_type_value);
}
if self.content_encoding != "" { if self.content_encoding != "" {
header.insert("Content-Encoding", HeaderValue::from_str(&self.content_encoding).expect("err")); if let Ok(encoding_value) = HeaderValue::from_str(&self.content_encoding) {
header.insert("Content-Encoding", encoding_value);
}
} }
if self.content_disposition != "" { if self.content_disposition != "" {
header.insert("Content-Disposition", HeaderValue::from_str(&self.content_disposition).expect("err")); if let Ok(disposition_value) = HeaderValue::from_str(&self.content_disposition) {
header.insert("Content-Disposition", disposition_value);
}
} }
if self.content_language != "" { if self.content_language != "" {
header.insert("Content-Language", HeaderValue::from_str(&self.content_language).expect("err")); if let Ok(language_value) = HeaderValue::from_str(&self.content_language) {
header.insert("Content-Language", language_value);
}
} }
if self.cache_control != "" { if self.cache_control != "" {
header.insert("Cache-Control", HeaderValue::from_str(&self.cache_control).expect("err")); if let Ok(cache_value) = HeaderValue::from_str(&self.cache_control) {
header.insert("Cache-Control", cache_value);
}
} }
if self.expires.unix_timestamp() != 0 { if self.expires.unix_timestamp() != 0 {
header.insert( if let Ok(expires_str) = self.expires.format(ISO8601_DATEFORMAT) {
"Expires", if let Ok(expires_value) = HeaderValue::from_str(&expires_str) {
HeaderValue::from_str(&self.expires.format(ISO8601_DATEFORMAT).unwrap()).expect("err"), header.insert("Expires", expires_value);
); //rustfs invalid header }
}
} }
if self.mode.as_str() != "" { if self.mode.as_str() != "" {
header.insert(X_AMZ_OBJECT_LOCK_MODE, HeaderValue::from_str(self.mode.as_str()).expect("err")); if let Ok(mode_value) = HeaderValue::from_str(self.mode.as_str()) {
header.insert(X_AMZ_OBJECT_LOCK_MODE, mode_value);
}
} }
if self.retain_until_date.unix_timestamp() != 0 { if self.retain_until_date.unix_timestamp() != 0 {
header.insert( if let Ok(retain_str) = self.retain_until_date.format(ISO8601_DATEFORMAT) {
X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, if let Ok(retain_value) = HeaderValue::from_str(&retain_str) {
HeaderValue::from_str(&self.retain_until_date.format(ISO8601_DATEFORMAT).unwrap()).expect("err"), header.insert(X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, retain_value);
); }
}
} }
if self.legalhold.as_str() != "" { if self.legalhold.as_str() != "" {
header.insert(X_AMZ_OBJECT_LOCK_LEGAL_HOLD, HeaderValue::from_str(self.legalhold.as_str()).expect("err")); if let Ok(legalhold_value) = HeaderValue::from_str(self.legalhold.as_str()) {
header.insert(X_AMZ_OBJECT_LOCK_LEGAL_HOLD, legalhold_value);
}
} }
if self.storage_class != "" { if self.storage_class != "" {
header.insert(X_AMZ_STORAGE_CLASS, HeaderValue::from_str(&self.storage_class).expect("err")); if let Ok(storage_class_value) = HeaderValue::from_str(&self.storage_class) {
header.insert(X_AMZ_STORAGE_CLASS, storage_class_value);
}
} }
if self.website_redirect_location != "" { if self.website_redirect_location != "" {
header.insert( if let Ok(redirect_value) = HeaderValue::from_str(&self.website_redirect_location) {
X_AMZ_WEBSITE_REDIRECT_LOCATION, header.insert(X_AMZ_WEBSITE_REDIRECT_LOCATION, redirect_value);
HeaderValue::from_str(&self.website_redirect_location).expect("err"), }
);
} }
if !self.internal.replication_status.as_str().is_empty() { if !self.internal.replication_status.as_str().is_empty() {
header.insert( if let Ok(replication_status_value) = HeaderValue::from_str(self.internal.replication_status.as_str()) {
X_AMZ_REPLICATION_STATUS, header.insert(X_AMZ_REPLICATION_STATUS, replication_status_value);
HeaderValue::from_str(self.internal.replication_status.as_str()).expect("err"), }
);
} }
for (k, v) in &self.user_metadata { for (k, v) in &self.user_metadata {
@@ -360,17 +376,21 @@ impl TransitionClient {
let mut md5_base64: String = "".to_string(); let mut md5_base64: String = "".to_string();
if opts.send_content_md5 { if opts.send_content_md5 {
let mut md5_hasher = self.md5_hasher.lock().unwrap(); if let Some(mut md5_hasher) = self.md5_hasher.lock().unwrap().as_mut() {
let hash = md5_hasher.as_mut().expect("err"); let hash = md5_hasher.hash_encode(&buf[..length]);
let hash = hash.hash_encode(&buf[..length]); md5_base64 = base64_encode(hash.as_ref());
md5_base64 = base64_encode(hash.as_ref()); }
} else { } else {
let mut crc = opts.auto_checksum.hasher()?; let mut crc = opts.auto_checksum.hasher()?;
crc.update(&buf[..length]); crc.update(&buf[..length]);
let csum = crc.finalize(); let csum = crc.finalize();
if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) { if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) {
custom_header.insert(header_name, base64_encode(csum.as_ref()).parse().expect("err")); if let Ok(header_value) = base64_encode(csum.as_ref()).parse() {
custom_header.insert(header_name, header_value);
} else {
warn!("Failed to parse checksum value");
}
} else { } else {
warn!("Invalid header name: {}", opts.auto_checksum.key()); warn!("Invalid header name: {}", opts.auto_checksum.key());
} }
@@ -127,7 +127,11 @@ impl TransitionClient {
let csum = crc.finalize(); let csum = crc.finalize();
if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) { if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) {
custom_header.insert(header_name, base64_encode(csum.as_ref()).parse().expect("err")); if let Ok(header_value) = base64_encode(csum.as_ref()).parse() {
custom_header.insert(header_name, header_value);
} else {
warn!("Failed to parse checksum value");
}
} else { } else {
warn!("Invalid header name: {}", opts.auto_checksum.key()); warn!("Invalid header name: {}", opts.auto_checksum.key());
} }
@@ -309,27 +313,27 @@ impl TransitionClient {
let h = resp.headers(); let h = resp.headers();
let mut obj_part = ObjectPart { let mut obj_part = ObjectPart {
checksum_crc32: if let Some(h_checksum_crc32) = h.get(ChecksumMode::ChecksumCRC32.key()) { checksum_crc32: if let Some(h_checksum_crc32) = h.get(ChecksumMode::ChecksumCRC32.key()) {
h_checksum_crc32.to_str().expect("err").to_string() h_checksum_crc32.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_crc32c: if let Some(h_checksum_crc32c) = h.get(ChecksumMode::ChecksumCRC32C.key()) { checksum_crc32c: if let Some(h_checksum_crc32c) = h.get(ChecksumMode::ChecksumCRC32C.key()) {
h_checksum_crc32c.to_str().expect("err").to_string() h_checksum_crc32c.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_sha1: if let Some(h_checksum_sha1) = h.get(ChecksumMode::ChecksumSHA1.key()) { checksum_sha1: if let Some(h_checksum_sha1) = h.get(ChecksumMode::ChecksumSHA1.key()) {
h_checksum_sha1.to_str().expect("err").to_string() h_checksum_sha1.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_sha256: if let Some(h_checksum_sha256) = h.get(ChecksumMode::ChecksumSHA256.key()) { checksum_sha256: if let Some(h_checksum_sha256) = h.get(ChecksumMode::ChecksumSHA256.key()) {
h_checksum_sha256.to_str().expect("err").to_string() h_checksum_sha256.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_crc64nvme: if let Some(h_checksum_crc64nvme) = h.get(ChecksumMode::ChecksumCRC64NVME.key()) { checksum_crc64nvme: if let Some(h_checksum_crc64nvme) = h.get(ChecksumMode::ChecksumCRC64NVME.key()) {
h_checksum_crc64nvme.to_str().expect("err").to_string() h_checksum_crc64nvme.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
@@ -338,7 +342,7 @@ impl TransitionClient {
obj_part.size = p.size; obj_part.size = p.size;
obj_part.part_num = p.part_number; obj_part.part_num = p.part_number;
obj_part.etag = if let Some(h_etag) = h.get("ETag") { obj_part.etag = if let Some(h_etag) = h.get("ETag") {
h_etag.to_str().expect("err").trim_matches('"').to_string() h_etag.to_str().unwrap_or("").trim_matches('"').to_string()
} else { } else {
"".to_string() "".to_string()
}; };
@@ -398,7 +402,7 @@ impl TransitionClient {
key: complete_multipart_upload_result.key, key: complete_multipart_upload_result.key,
etag: trim_etag(&complete_multipart_upload_result.etag), etag: trim_etag(&complete_multipart_upload_result.etag),
version_id: if let Some(h_x_amz_version_id) = h.get(X_AMZ_VERSION_ID) { version_id: if let Some(h_x_amz_version_id) = h.get(X_AMZ_VERSION_ID) {
h_x_amz_version_id.to_str().expect("err").to_string() h_x_amz_version_id.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
@@ -21,6 +21,7 @@
use bytes::Bytes; use bytes::Bytes;
use futures::future::join_all; use futures::future::join_all;
use http::{HeaderMap, HeaderName, HeaderValue, StatusCode}; use http::{HeaderMap, HeaderName, HeaderValue, StatusCode};
use std::io::Error;
use std::sync::RwLock; use std::sync::RwLock;
use std::{collections::HashMap, sync::Arc}; use std::{collections::HashMap, sync::Arc};
use time::{OffsetDateTime, format_description}; use time::{OffsetDateTime, format_description};
@@ -152,7 +153,10 @@ impl TransitionClient {
if opts.send_content_md5 { if opts.send_content_md5 {
let mut md5_hasher = self.md5_hasher.lock().unwrap(); let mut md5_hasher = self.md5_hasher.lock().unwrap();
let md5_hash = md5_hasher.as_mut().expect("err"); let md5_hash = match md5_hasher.as_mut() {
Some(hasher) => hasher,
None => return Err(std::io::Error::other("MD5 hasher not initialized")),
};
let hash = md5_hash.hash_encode(&buf[..length]); let hash = md5_hash.hash_encode(&buf[..length]);
md5_base64 = base64_encode(hash.as_ref()); md5_base64 = base64_encode(hash.as_ref());
} else { } else {
@@ -161,7 +165,11 @@ impl TransitionClient {
let csum = crc.finalize(); let csum = crc.finalize();
if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) { if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) {
custom_header.insert(header_name, base64_encode(csum.as_ref()).parse().expect("err")); if let Ok(header_value) = base64_encode(csum.as_ref()).parse() {
custom_header.insert(header_name, header_value);
} else {
warn!("Failed to parse checksum value");
}
} else { } else {
warn!("Invalid header name: {}", opts.auto_checksum.key()); warn!("Invalid header name: {}", opts.auto_checksum.key());
} }
@@ -275,11 +283,14 @@ impl TransitionClient {
for part_number in 1..=total_parts_count { for part_number in 1..=total_parts_count {
let mut buf = Vec::<u8>::new(); let mut buf = Vec::<u8>::new();
select! { select! {
buf = bufs_rx.recv() => {} buf1 = bufs_rx.recv() => {
if let Some(buf1) = buf1 {
buf = buf1;
}
}
err = err_rx.recv() => { err = err_rx.recv() => {
//cancel_token.cancel(); //cancel_token.cancel();
//wg.Wait() return Err(err.unwrap_or_else(|| std::io::Error::other("Unknown error received from channel")));
return Err(err.expect("err"));
} }
else => (), else => (),
} }
@@ -309,7 +320,11 @@ impl TransitionClient {
let csum = crc.finalize(); let csum = crc.finalize();
if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) { if let Ok(header_name) = HeaderName::from_bytes(opts.auto_checksum.key().as_bytes()) {
custom_header.insert(header_name, base64_encode(csum.as_ref()).parse().expect("err")); if let Ok(header_value) = base64_encode(csum.as_ref()).parse() {
custom_header.insert(header_name, header_value);
} else {
warn!("Failed to parse checksum value");
}
} else { } else {
warn!("Invalid header name: {}", opts.auto_checksum.key()); warn!("Invalid header name: {}", opts.auto_checksum.key());
} }
@@ -319,12 +334,19 @@ impl TransitionClient {
let clone_parts_info = parts_info.clone(); let clone_parts_info = parts_info.clone();
let clone_upload_id = upload_id.clone(); let clone_upload_id = upload_id.clone();
let clone_self = self.clone(); let clone_self = self.clone();
let err_tx_clone = err_tx.clone();
futures.push(async move { futures.push(async move {
let mut md5_base64: String = "".to_string(); let mut md5_base64: String = "".to_string();
if opts.send_content_md5 { if opts.send_content_md5 {
let mut md5_hasher = clone_self.md5_hasher.lock().unwrap(); let mut md5_hasher = clone_self.md5_hasher.lock().unwrap();
let md5_hash = md5_hasher.as_mut().expect("err"); let md5_hash = match md5_hasher.as_mut() {
Some(hasher) => hasher,
None => {
//let _ = err_tx_clone.send(std::io::Error::other("MD5 hasher not initialized")).await;
return Ok::<(), Error>(());
}
};
let hash = md5_hash.hash_encode(&buf[..length]); let hash = md5_hash.hash_encode(&buf[..length]);
md5_base64 = base64_encode(hash.as_ref()); md5_base64 = base64_encode(hash.as_ref());
} }
@@ -344,12 +366,21 @@ impl TransitionClient {
sha256_hex: "".to_string(), sha256_hex: "".to_string(),
trailer: HeaderMap::new(), trailer: HeaderMap::new(),
}; };
let obj_part = clone_self.upload_part(&mut p).await.expect("err"); let obj_part = match clone_self.upload_part(&mut p).await {
Ok(part) => part,
Err(err) => {
let _ = err_tx_clone.send(std::io::Error::other(err.to_string())).await;
return Err::<(), Error>(err);
}
};
let mut clone_parts_info = clone_parts_info.write().unwrap(); {
clone_parts_info.entry(part_number).or_insert(obj_part); let mut clone_parts_info = clone_parts_info.write().unwrap();
clone_parts_info.entry(part_number).or_insert(obj_part);
}
clone_bufs_tx.send(buf); let _ = clone_bufs_tx.send(buf).await;
Ok::<(), Error>(())
}); });
total_uploaded_size += length as i64; total_uploaded_size += length as i64;
@@ -359,7 +390,7 @@ impl TransitionClient {
select! { select! {
err = err_rx.recv() => { err = err_rx.recv() => {
return Err(err.expect("err")); return Err(err.unwrap_or_else(|| std::io::Error::other("Unknown error received from channel")));
} }
else => (), else => (),
} }
@@ -504,9 +535,10 @@ impl TransitionClient {
Ok(UploadInfo { Ok(UploadInfo {
bucket: bucket_name.to_string(), bucket: bucket_name.to_string(),
key: object_name.to_string(), key: object_name.to_string(),
etag: trim_etag(h.get("ETag").expect("err").to_str().expect("err")), etag: trim_etag(h.get("ETag").and_then(|v| v.to_str().ok()).unwrap_or("")),
version_id: if let Some(h_x_amz_version_id) = h.get(X_AMZ_VERSION_ID) { version_id: if let Some(h_x_amz_version_id) = h.get(X_AMZ_VERSION_ID) {
h_x_amz_version_id.to_str().expect("err").to_string() h_x_amz_version_id.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
@@ -514,27 +546,27 @@ impl TransitionClient {
expiration: exp_time, expiration: exp_time,
expiration_rule_id: rule_id, expiration_rule_id: rule_id,
checksum_crc32: if let Some(h_checksum_crc32) = h.get(ChecksumMode::ChecksumCRC32.key()) { checksum_crc32: if let Some(h_checksum_crc32) = h.get(ChecksumMode::ChecksumCRC32.key()) {
h_checksum_crc32.to_str().expect("err").to_string() h_checksum_crc32.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_crc32c: if let Some(h_checksum_crc32c) = h.get(ChecksumMode::ChecksumCRC32C.key()) { checksum_crc32c: if let Some(h_checksum_crc32c) = h.get(ChecksumMode::ChecksumCRC32C.key()) {
h_checksum_crc32c.to_str().expect("err").to_string() h_checksum_crc32c.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_sha1: if let Some(h_checksum_sha1) = h.get(ChecksumMode::ChecksumSHA1.key()) { checksum_sha1: if let Some(h_checksum_sha1) = h.get(ChecksumMode::ChecksumSHA1.key()) {
h_checksum_sha1.to_str().expect("err").to_string() h_checksum_sha1.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_sha256: if let Some(h_checksum_sha256) = h.get(ChecksumMode::ChecksumSHA256.key()) { checksum_sha256: if let Some(h_checksum_sha256) = h.get(ChecksumMode::ChecksumSHA256.key()) {
h_checksum_sha256.to_str().expect("err").to_string() h_checksum_sha256.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
checksum_crc64nvme: if let Some(h_checksum_crc64nvme) = h.get(ChecksumMode::ChecksumCRC64NVME.key()) { checksum_crc64nvme: if let Some(h_checksum_crc64nvme) = h.get(ChecksumMode::ChecksumCRC64NVME.key()) {
h_checksum_crc64nvme.to_str().expect("err").to_string() h_checksum_crc64nvme.to_str().unwrap_or("").to_string()
} else { } else {
"".to_string() "".to_string()
}, },
+26 -27
View File
@@ -25,7 +25,7 @@ use hyper::body::Bytes;
use rustfs_utils::HashAlgorithm; use rustfs_utils::HashAlgorithm;
use s3s::S3ErrorCode; use s3s::S3ErrorCode;
use s3s::dto::ReplicationStatus; use s3s::dto::ReplicationStatus;
use s3s::header::X_AMZ_BYPASS_GOVERNANCE_RETENTION; use s3s::header::{X_AMZ_BYPASS_GOVERNANCE_RETENTION, X_AMZ_DELETE_MARKER, X_AMZ_VERSION_ID};
use serde::Deserialize; use serde::Deserialize;
use std::fmt::Display; use std::fmt::Display;
use std::{ use std::{
@@ -111,8 +111,9 @@ impl TransitionClient {
.await?; .await?;
{ {
let mut bucket_loc_cache = self.bucket_loc_cache.lock().unwrap(); if let Ok(mut bucket_loc_cache) = self.bucket_loc_cache.lock() {
bucket_loc_cache.delete(bucket_name); bucket_loc_cache.delete(bucket_name);
}
} }
Ok(()) Ok(())
} }
@@ -142,8 +143,9 @@ impl TransitionClient {
.await?; .await?;
{ {
let mut bucket_loc_cache = self.bucket_loc_cache.lock().unwrap(); if let Ok(mut bucket_loc_cache) = self.bucket_loc_cache.lock() {
bucket_loc_cache.delete(bucket_name); bucket_loc_cache.delete(bucket_name);
}
} }
Ok(()) Ok(())
@@ -168,7 +170,7 @@ impl TransitionClient {
let mut headers = HeaderMap::new(); let mut headers = HeaderMap::new();
if opts.governance_bypass { if opts.governance_bypass {
headers.insert(X_AMZ_BYPASS_GOVERNANCE_RETENTION, "true".parse().expect("err")); //amzBypassGovernance headers.insert(X_AMZ_BYPASS_GOVERNANCE_RETENTION, HeaderValue::from_static("true")); //amzBypassGovernance
} }
let resp = self let resp = self
@@ -197,13 +199,12 @@ impl TransitionClient {
Ok(RemoveObjectResult { Ok(RemoveObjectResult {
object_name: object_name.to_string(), object_name: object_name.to_string(),
object_version_id: opts.version_id, object_version_id: opts.version_id,
delete_marker: resp.headers().get("x-amz-delete-marker").expect("err") == "true", delete_marker: resp.headers().get(X_AMZ_DELETE_MARKER).map_or(false, |v| v == "true"),
delete_marker_version_id: resp delete_marker_version_id: resp
.headers() .headers()
.get("x-amz-version-id") .get(X_AMZ_VERSION_ID)
.expect("err") .and_then(|v| v.to_str().ok())
.to_str() .unwrap_or_default()
.expect("err")
.to_string(), .to_string(),
..Default::default() ..Default::default()
}) })
@@ -290,15 +291,15 @@ impl TransitionClient {
bucket_name, bucket_name,
&object.name, &object.name,
RemoveObjectOptions { RemoveObjectOptions {
version_id: object.version_id.expect("err").to_string(), version_id: object.version_id.map(|id| id.to_string()).unwrap_or_default(),
governance_bypass: opts.governance_bypass, governance_bypass: opts.governance_bypass,
..Default::default() ..Default::default()
}, },
) )
.await?; .await?;
let remove_result_clone = remove_result.clone(); let remove_result_clone = remove_result.clone();
if !remove_result.err.is_none() { if let Some(err) = &remove_result.err {
match to_error_response(&remove_result.err.expect("err")).code { match to_error_response(err).code {
S3ErrorCode::InvalidArgument | S3ErrorCode::NoSuchVersion => { S3ErrorCode::InvalidArgument | S3ErrorCode::NoSuchVersion => {
continue; continue;
} }
@@ -326,7 +327,7 @@ impl TransitionClient {
let mut headers = HeaderMap::new(); let mut headers = HeaderMap::new();
if opts.governance_bypass { if opts.governance_bypass {
headers.insert(X_AMZ_BYPASS_GOVERNANCE_RETENTION, "true".parse().expect("err")); headers.insert(X_AMZ_BYPASS_GOVERNANCE_RETENTION, HeaderValue::from_static("true"));
} }
let remove_bytes = generate_remove_multi_objects_request(&batch); let remove_bytes = generate_remove_multi_objects_request(&batch);
@@ -423,23 +424,20 @@ impl TransitionClient {
request_id: resp request_id: resp
.headers() .headers()
.get("x-amz-request-id") .get("x-amz-request-id")
.expect("err") .and_then(|v| v.to_str().ok())
.to_str() .unwrap_or_default()
.expect("err")
.to_string(), .to_string(),
host_id: resp host_id: resp
.headers() .headers()
.get("x-amz-id-2") .get("x-amz-id-2")
.expect("err") .and_then(|v| v.to_str().ok())
.to_str() .unwrap_or_default()
.expect("err")
.to_string(), .to_string(),
region: resp region: resp
.headers() .headers()
.get("x-amz-bucket-region") .get("x-amz-bucket-region")
.expect("err") .and_then(|v| v.to_str().ok())
.to_str() .unwrap_or_default()
.expect("err")
.to_string(), .to_string(),
..Default::default() ..Default::default()
}; };
@@ -472,10 +470,11 @@ pub struct RemoveObjectError {
impl Display for RemoveObjectError { impl Display for RemoveObjectError {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result { fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
if self.err.is_none() { if let Some(err) = &self.err {
return write!(f, "unexpected remove object error result"); write!(f, "{}", err.to_string())
} else {
write!(f, "unexpected remove object error result")
} }
write!(f, "{}", self.err.as_ref().expect("err").to_string())
} }
} }
+34 -22
View File
@@ -70,10 +70,10 @@ impl TransitionClient {
let mut location; let mut location;
{ {
let mut bucket_loc_cache = self.bucket_loc_cache.lock().unwrap(); if let Ok(bucket_loc_cache) = self.bucket_loc_cache.lock() {
let ret = bucket_loc_cache.get(bucket_name); if let Some(location) = bucket_loc_cache.get(bucket_name) {
if let Some(location) = ret { return Ok(location);
return Ok(location); }
} }
//location = ret?; //location = ret?;
} }
@@ -83,8 +83,9 @@ impl TransitionClient {
let mut resp = self.doit(req).await?; let mut resp = self.doit(req).await?;
location = process_bucket_location_response(resp, bucket_name, &self.tier_type).await?; location = process_bucket_location_response(resp, bucket_name, &self.tier_type).await?;
{ {
let mut bucket_loc_cache = self.bucket_loc_cache.lock().unwrap(); if let Ok(mut bucket_loc_cache) = self.bucket_loc_cache.lock() {
bucket_loc_cache.set(bucket_name, &location); bucket_loc_cache.set(bucket_name, &location);
}
} }
Ok(location) Ok(location)
} }
@@ -108,7 +109,11 @@ impl TransitionClient {
url_str.push_str("://"); url_str.push_str("://");
url_str.push_str(bucket_name); url_str.push_str(bucket_name);
url_str.push_str("."); url_str.push_str(".");
url_str.push_str(target_url.host_str().expect("err")); url_str.push_str(
target_url
.host_str()
.ok_or_else(|| std::io::Error::new(std::io::ErrorKind::InvalidInput, "host is none"))?,
);
url_str.push_str("/?location"); url_str.push_str("/?location");
} else { } else {
let mut path = bucket_name.to_string(); let mut path = bucket_name.to_string();
@@ -135,13 +140,16 @@ impl TransitionClient {
let value; let value;
{ {
let mut creds_provider = self.creds_provider.lock().unwrap(); if let Ok(mut creds_provider) = self.creds_provider.lock() {
value = match creds_provider.get_with_context(Some(self.cred_context())) { value = match creds_provider.get_with_context(Some(self.cred_context())) {
Ok(v) => v, Ok(v) => v,
Err(err) => { Err(err) => {
return Err(std::io::Error::other(err)); return Err(std::io::Error::other(err));
} }
}; };
} else {
return Err(std::io::Error::other("Failed to acquire credentials provider lock"));
}
} }
let mut signer_type = value.signer_type.clone(); let mut signer_type = value.signer_type.clone();
@@ -171,8 +179,9 @@ impl TransitionClient {
content_sha256 = UNSIGNED_PAYLOAD.to_string(); content_sha256 = UNSIGNED_PAYLOAD.to_string();
} }
req.headers_mut() if let Ok(content_sha256_value) = content_sha256.parse() {
.insert("X-Amz-Content-Sha256", content_sha256.parse().unwrap()); req.headers_mut().insert("X-Amz-Content-Sha256", content_sha256_value);
}
let req = rustfs_signer::sign_v4(req, 0, &access_key_id, &secret_access_key, &session_token, "us-east-1"); let req = rustfs_signer::sign_v4(req, 0, &access_key_id, &secret_access_key, &session_token, "us-east-1");
Ok(req) Ok(req)
} }
@@ -228,13 +237,16 @@ async fn process_bucket_location_response(
} }
let mut location = "".to_string(); let mut location = "".to_string();
if tier_type == "huaweicloud" { if tier_type == "huaweicloud" {
let d = quick_xml::de::from_str::<CreateBucketConfiguration>(&String::from_utf8(body_vec).unwrap()).unwrap(); if let Ok(body_str) = String::from_utf8(body_vec) {
location = d.location_constraint; if let Ok(d) = quick_xml::de::from_str::<CreateBucketConfiguration>(&body_str) {
location = d.location_constraint;
}
}
} else { } else {
if let Ok(LocationConstraint { field }) = if let Ok(body_str) = String::from_utf8(body_vec) {
quick_xml::de::from_str::<LocationConstraint>(&String::from_utf8(body_vec).unwrap()) if let Ok(LocationConstraint { field }) = quick_xml::de::from_str::<LocationConstraint>(&body_str) {
{ location = field;
location = field; }
} }
} }
//debug!("location: {}", location); //debug!("location: {}", location);
+4 -2
View File
@@ -1,4 +1,3 @@
#![allow(unused_imports)]
// Copyright 2024 RustFS Team // Copyright 2024 RustFS Team
// //
// Licensed under the Apache License, Version 2.0 (the "License"); // Licensed under the Apache License, Version 2.0 (the "License");
@@ -12,6 +11,7 @@
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
#![allow(unused_imports)]
#![allow(unused_variables)] #![allow(unused_variables)]
#![allow(unused_mut)] #![allow(unused_mut)]
#![allow(unused_assignments)] #![allow(unused_assignments)]
@@ -57,7 +57,9 @@ impl<P: Provider + Default> Credentials<P> {
pub fn get_with_context(&mut self, mut cc: Option<CredContext>) -> Result<Value, std::io::Error> { pub fn get_with_context(&mut self, mut cc: Option<CredContext>) -> Result<Value, std::io::Error> {
if self.is_expired() { if self.is_expired() {
let creds = self.provider.retrieve_with_cred_context(cc.expect("err")); let creds = self.provider.retrieve_with_cred_context(cc.unwrap_or(CredContext {
endpoint: "".to_string(),
}));
self.creds = creds; self.creds = creds;
self.force_refresh = false; self.force_refresh = false;
} }
+64 -36
View File
@@ -251,9 +251,10 @@ impl TransitionClient {
}; };
{ {
let mut md5_hasher = client.md5_hasher.lock().unwrap(); if let Ok(mut md5_hasher) = client.md5_hasher.lock() {
if md5_hasher.is_none() { if md5_hasher.is_none() {
*md5_hasher = Some(HashAlgorithm::Md5); *md5_hasher = Some(HashAlgorithm::Md5);
}
} }
} }
if client.sha256_hasher.is_none() { if client.sha256_hasher.is_none() {
@@ -275,25 +276,30 @@ impl TransitionClient {
} }
fn trace_errors_only_off(&self) { fn trace_errors_only_off(&self) {
let mut trace_errors_only = self.trace_errors_only.lock().unwrap(); if let Ok(mut trace_errors_only) = self.trace_errors_only.lock() {
*trace_errors_only = false; *trace_errors_only = false;
}
} }
fn trace_off(&self) { fn trace_off(&self) {
let mut is_trace_enabled = self.is_trace_enabled.lock().unwrap(); if let Ok(mut is_trace_enabled) = self.is_trace_enabled.lock() {
*is_trace_enabled = false; *is_trace_enabled = false;
let mut trace_errors_only = self.trace_errors_only.lock().unwrap(); }
*trace_errors_only = false; if let Ok(mut trace_errors_only) = self.trace_errors_only.lock() {
*trace_errors_only = false;
}
} }
fn set_s3_transfer_accelerate(&self, accelerate_endpoint: &str) { fn set_s3_transfer_accelerate(&self, accelerate_endpoint: &str) {
let mut endpoint = self.s3_accelerate_endpoint.lock().unwrap(); if let Ok(mut endpoint) = self.s3_accelerate_endpoint.lock() {
*endpoint = accelerate_endpoint.to_string(); *endpoint = accelerate_endpoint.to_string();
}
} }
fn set_s3_enable_dual_stack(&self, enabled: bool) { fn set_s3_enable_dual_stack(&self, enabled: bool) {
let mut dual_stack = self.s3_dual_stack_enabled.lock().unwrap(); if let Ok(mut dual_stack) = self.s3_dual_stack_enabled.lock() {
*dual_stack = enabled; *dual_stack = enabled;
}
} }
pub fn hash_materials( pub fn hash_materials(
@@ -352,7 +358,6 @@ impl TransitionClient {
let resp; let resp;
let http_client = self.http_client.clone(); let http_client = self.http_client.clone();
{ {
//let mut http_client = http_client.lock().unwrap();
req_method = req.method().clone(); req_method = req.method().clone();
req_uri = req.uri().clone(); req_uri = req.uri().clone();
req_headers = req.headers().clone(); req_headers = req.headers().clone();
@@ -368,7 +373,10 @@ impl TransitionClient {
return Err(std::io::Error::other(err)); return Err(std::io::Error::other(err));
} }
let resp = resp.unwrap(); let resp = match resp {
Ok(r) => r,
Err(_) => return Err(std::io::Error::other("Unexpected error in response")),
};
debug!("http_resp: {:?}", resp); debug!("http_resp: {:?}", resp);
//let b = resp.body_mut().store_all_unlimited().await.unwrap().to_vec(); //let b = resp.body_mut().store_all_unlimited().await.unwrap().to_vec();
@@ -455,11 +463,13 @@ impl TransitionClient {
return Err(std::io::Error::other(err_response)); return Err(std::io::Error::other(err_response));
} }
if metadata.bucket_name != "" { if metadata.bucket_name != "" {
let mut bucket_loc_cache = self.bucket_loc_cache.lock().unwrap(); if let Ok(mut bucket_loc_cache) = self.bucket_loc_cache.lock() {
let location = bucket_loc_cache.get(&metadata.bucket_name); if let Some(location) = bucket_loc_cache.get(&metadata.bucket_name) {
if location.is_some() && location.unwrap() != err_response.region { if location != err_response.region {
bucket_loc_cache.set(&metadata.bucket_name, &err_response.region); bucket_loc_cache.set(&metadata.bucket_name, &err_response.region);
//continue; //continue;
}
}
} }
} else if err_response.region != metadata.bucket_location { } else if err_response.region != metadata.bucket_location {
metadata.bucket_location = err_response.region.clone(); metadata.bucket_location = err_response.region.clone();
@@ -518,8 +528,11 @@ impl TransitionClient {
let value; let value;
{ {
let mut creds_provider = self.creds_provider.lock().unwrap(); if let Ok(mut creds_provider) = self.creds_provider.lock() {
value = creds_provider.get_with_context(Some(self.cred_context()))?; value = creds_provider.get_with_context(Some(self.cred_context()))?;
} else {
return Err(std::io::Error::other("Failed to acquire credentials provider lock"));
}
} }
let mut signer_type = value.signer_type.clone(); let mut signer_type = value.signer_type.clone();
@@ -548,8 +561,10 @@ impl TransitionClient {
))); )));
} }
let headers = req.headers_mut(); let headers = req.headers_mut();
for (k, v) in metadata.extra_pre_sign_header.as_ref().unwrap() { if let Some(extra_headers) = metadata.extra_pre_sign_header.as_ref() {
headers.insert(k, v.clone()); for (k, v) in extra_headers {
headers.insert(k, v.clone());
}
} }
} }
if signer_type == SignatureType::SignatureV2 { if signer_type == SignatureType::SignatureV2 {
@@ -571,18 +586,22 @@ impl TransitionClient {
self.set_user_agent(&mut req); self.set_user_agent(&mut req);
for (k, v) in metadata.custom_header.clone() { for (k, v) in metadata.custom_header.clone() {
req.headers_mut().insert(k.expect("err"), v); if let Some(key) = k {
req.headers_mut().insert(key, v);
}
} }
//req.content_length = metadata.content_length; //req.content_length = metadata.content_length;
if metadata.content_length <= -1 { if metadata.content_length <= -1 {
let chunked_value = HeaderValue::from_str(&vec!["chunked"].join(",")).expect("err"); if let Ok(chunked_value) = HeaderValue::from_str(&vec!["chunked"].join(",")) {
req.headers_mut().insert(http::header::TRANSFER_ENCODING, chunked_value); req.headers_mut().insert(http::header::TRANSFER_ENCODING, chunked_value);
}
} }
if metadata.content_md5_base64.len() > 0 { if metadata.content_md5_base64.len() > 0 {
let md5_value = HeaderValue::from_str(&metadata.content_md5_base64).expect("err"); if let Ok(md5_value) = HeaderValue::from_str(&metadata.content_md5_base64) {
req.headers_mut().insert("Content-Md5", md5_value); req.headers_mut().insert("Content-Md5", md5_value);
}
} }
if signer_type == SignatureType::SignatureAnonymous { if signer_type == SignatureType::SignatureAnonymous {
@@ -607,8 +626,13 @@ impl TransitionClient {
} else if metadata.trailer.len() > 0 { } else if metadata.trailer.len() > 0 {
sha_header = UNSIGNED_PAYLOAD_TRAILER.to_string(); sha_header = UNSIGNED_PAYLOAD_TRAILER.to_string();
} }
req.headers_mut() let header_name = "X-Amz-Content-Sha256"
.insert("X-Amz-Content-Sha256".parse::<HeaderName>().unwrap(), sha_header.parse().expect("err")); .parse::<HeaderName>()
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let header_value = sha_header
.parse()
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
req.headers_mut().insert(header_name, header_value);
req = rustfs_signer::sign_v4_trailer( req = rustfs_signer::sign_v4_trailer(
req, req,
@@ -636,7 +660,7 @@ impl TransitionClient {
pub fn set_user_agent(&self, req: &mut Request<s3s::Body>) { pub fn set_user_agent(&self, req: &mut Request<s3s::Body>) {
let headers = req.headers_mut(); let headers = req.headers_mut();
headers.insert("User-Agent", C_USER_AGENT.parse().expect("err")); headers.insert("User-Agent", HeaderValue::from_static(C_USER_AGENT));
} }
fn make_target_url( fn make_target_url(
@@ -648,7 +672,10 @@ impl TransitionClient {
query_values: &HashMap<String, String>, query_values: &HashMap<String, String>,
) -> Result<Url, std::io::Error> { ) -> Result<Url, std::io::Error> {
let scheme = self.endpoint_url.scheme(); let scheme = self.endpoint_url.scheme();
let host = self.endpoint_url.host().unwrap(); let host = self
.endpoint_url
.host()
.ok_or_else(|| std::io::Error::other("Endpoint URL has no host"))?;
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let port = self.endpoint_url.port().unwrap_or(default_port); let port = self.endpoint_url.port().unwrap_or(default_port);
@@ -1155,9 +1182,10 @@ pub fn to_object_info(bucket_name: &str, object_name: &str, h: &HeaderMap) -> Re
for (name, value) in h.iter() { for (name, value) in h.iter() {
let header_name = name.as_str().to_lowercase(); let header_name = name.as_str().to_lowercase();
if header_name.starts_with("x-amz-meta-") { if header_name.starts_with("x-amz-meta-") {
let key = header_name.strip_prefix("x-amz-meta-").unwrap().to_string(); if let Some(key) = header_name.strip_prefix("x-amz-meta-") {
if let Ok(value_str) = value.to_str() { if let Ok(value_str) = value.to_str() {
meta.insert(key, value_str.to_string()); meta.insert(key.to_string(), value_str.to_string());
}
} }
} }
} }
+87 -67
View File
@@ -808,14 +808,19 @@ impl TierConfigMgr {
} }
} }
if !force { if !force {
let inuse = d.expect("err").in_use().await; if let Ok(driver) = d {
if let Err(err) = inuse { match driver.in_use().await {
let mut e = ERR_TIER_PERM_ERR.clone(); Err(err) => {
e.message.push('.'); let mut e = ERR_TIER_PERM_ERR.clone();
e.message.push_str(&err.to_string()); e.message.push('.');
return Err(e); e.message.push_str(&err.to_string());
} else if inuse.expect("err") { return Err(e);
return Err(ERR_TIER_BACKEND_NOT_EMPTY.clone()); }
Ok(in_use) if in_use => {
return Err(ERR_TIER_BACKEND_NOT_EMPTY.clone());
}
_ => {}
}
} }
} }
self.tiers.remove(tier_name); self.tiers.remove(tier_name);
@@ -842,11 +847,11 @@ impl TierConfigMgr {
} }
pub fn tier_type(&self, tier_name: &str) -> String { pub fn tier_type(&self, tier_name: &str) -> String {
let cfg = self.tiers.get(tier_name); if let Some(cfg) = self.tiers.get(tier_name) {
if cfg.is_none() { cfg.tier_type.as_lowercase()
return "internal".to_string(); } else {
"internal".to_string()
} }
cfg.expect("err").tier_type.as_lowercase()
} }
pub fn list_tiers(&self) -> Vec<TierConfig> { pub fn list_tiers(&self) -> Vec<TierConfig> {
@@ -876,81 +881,90 @@ impl TierConfigMgr {
let mut tier_config = self.tiers[tier_name].clone(); let mut tier_config = self.tiers[tier_name].clone();
match tier_type { match tier_type {
TierType::S3 => { TierType::S3 => {
let mut s3 = tier_config.s3.as_mut().expect("err"); if let Some(s3) = tier_config.s3.as_mut() {
if creds.aws_role { if creds.aws_role {
s3.aws_role = true s3.aws_role = true
} }
if creds.aws_role_web_identity_token_file != "" && creds.aws_role_arn != "" { if creds.aws_role_web_identity_token_file != "" && creds.aws_role_arn != "" {
s3.aws_role_arn = creds.aws_role_arn; s3.aws_role_arn = creds.aws_role_arn;
s3.aws_role_web_identity_token_file = creds.aws_role_web_identity_token_file; s3.aws_role_web_identity_token_file = creds.aws_role_web_identity_token_file;
} }
if creds.access_key != "" && creds.secret_key != "" { if creds.access_key != "" && creds.secret_key != "" {
s3.access_key = creds.access_key; s3.access_key = creds.access_key;
s3.secret_key = creds.secret_key; s3.secret_key = creds.secret_key;
}
} }
} }
TierType::RustFS => { TierType::RustFS => {
let mut rustfs = tier_config.rustfs.as_mut().expect("err"); if let Some(rustfs) = tier_config.rustfs.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
rustfs.access_key = creds.access_key;
rustfs.secret_key = creds.secret_key;
} }
rustfs.access_key = creds.access_key;
rustfs.secret_key = creds.secret_key;
} }
TierType::MinIO => { TierType::MinIO => {
let compatible_backend = tier_config.minio.as_mut().expect("err"); if let Some(compatible_backend) = tier_config.minio.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
compatible_backend.access_key = creds.access_key;
compatible_backend.secret_key = creds.secret_key;
} }
compatible_backend.access_key = creds.access_key;
compatible_backend.secret_key = creds.secret_key;
} }
TierType::Aliyun => { TierType::Aliyun => {
let mut aliyun = tier_config.aliyun.as_mut().expect("err"); if let Some(aliyun) = tier_config.aliyun.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
aliyun.access_key = creds.access_key;
aliyun.secret_key = creds.secret_key;
} }
aliyun.access_key = creds.access_key;
aliyun.secret_key = creds.secret_key;
} }
TierType::Tencent => { TierType::Tencent => {
let mut tencent = tier_config.tencent.as_mut().expect("err"); if let Some(tencent) = tier_config.tencent.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
tencent.access_key = creds.access_key;
tencent.secret_key = creds.secret_key;
} }
tencent.access_key = creds.access_key;
tencent.secret_key = creds.secret_key;
} }
TierType::Huaweicloud => { TierType::Huaweicloud => {
let mut huaweicloud = tier_config.huaweicloud.as_mut().expect("err"); if let Some(huaweicloud) = tier_config.huaweicloud.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
huaweicloud.access_key = creds.access_key;
huaweicloud.secret_key = creds.secret_key;
} }
huaweicloud.access_key = creds.access_key;
huaweicloud.secret_key = creds.secret_key;
} }
TierType::Azure => { TierType::Azure => {
let mut azure = tier_config.azure.as_mut().expect("err"); if let Some(azure) = tier_config.azure.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
azure.access_key = creds.access_key;
azure.secret_key = creds.secret_key;
} }
azure.access_key = creds.access_key;
azure.secret_key = creds.secret_key;
} }
TierType::GCS => { TierType::GCS => {
let mut gcs = tier_config.gcs.as_mut().expect("err"); if let Some(gcs) = tier_config.gcs.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
gcs.creds = creds.access_key; //creds.creds_json
} }
gcs.creds = creds.access_key; //creds.creds_json
} }
TierType::R2 => { TierType::R2 => {
let mut r2 = tier_config.r2.as_mut().expect("err"); if let Some(r2) = tier_config.r2.as_mut() {
if creds.access_key == "" || creds.secret_key == "" { if creds.access_key == "" || creds.secret_key == "" {
return Err(ERR_TIER_MISSING_CREDENTIALS.clone()); return Err(ERR_TIER_MISSING_CREDENTIALS.clone());
}
r2.access_key = creds.access_key;
r2.secret_key = creds.secret_key;
} }
r2.access_key = creds.access_key;
r2.secret_key = creds.secret_key;
} }
_ => (), _ => (),
} }
@@ -964,7 +978,7 @@ impl TierConfigMgr {
pub async fn get_driver<'a>(&'a mut self, tier_name: &str) -> std::result::Result<&'a WarmBackendImpl, AdminError> { pub async fn get_driver<'a>(&'a mut self, tier_name: &str) -> std::result::Result<&'a WarmBackendImpl, AdminError> {
// Return cached driver if present // Return cached driver if present
if self.driver_cache.contains_key(tier_name) { if self.driver_cache.contains_key(tier_name) {
return Ok(self.driver_cache.get(tier_name).unwrap()); return Ok(self.driver_cache.get(tier_name).expect("Driver not found in cache"));
} }
// Get tier configuration and create new driver // Get tier configuration and create new driver
@@ -974,7 +988,10 @@ impl TierConfigMgr {
// Insert and return reference // Insert and return reference
self.driver_cache.insert(tier_name.to_string(), driver); self.driver_cache.insert(tier_name.to_string(), driver);
Ok(self.driver_cache.get(tier_name).unwrap()) Ok(self
.driver_cache
.get(tier_name)
.expect("Driver not found in cache after insertion"))
} }
pub async fn reload(&mut self, api: Arc<ECStore>) -> std::result::Result<(), std::io::Error> { pub async fn reload(&mut self, api: Arc<ECStore>) -> std::result::Result<(), std::io::Error> {
@@ -989,9 +1006,12 @@ impl TierConfigMgr {
} }
self.driver_cache.clear(); self.driver_cache.clear();
self.tiers.clear(); self.tiers.clear();
let new_config = new_config.expect("err"); if let Ok(config) = new_config {
for (tier, cfg) in new_config.tiers { for (tier, cfg) in config.tiers {
self.tiers.insert(tier, cfg); self.tiers.insert(tier, cfg);
}
} else {
return Err(std::io::Error::other("Failed to load tier configuration"));
} }
self.last_refreshed_at = OffsetDateTime::now_utc(); self.last_refreshed_at = OffsetDateTime::now_utc();
Ok(()) Ok(())
+81 -63
View File
@@ -155,49 +155,67 @@ impl Clone for TierConfig {
let mut r2 = None; let mut r2 = None;
match self.tier_type { match self.tier_type {
TierType::S3 => { TierType::S3 => {
let mut s3_ = self.s3.as_ref().expect("err").clone(); if let Some(s3_) = self.s3.as_ref() {
s3_.secret_key = "REDACTED".to_string(); let mut s3_clone = s3_.clone();
s3 = Some(s3_); s3_clone.secret_key = "REDACTED".to_string();
s3 = Some(s3_clone);
}
} }
TierType::RustFS => { TierType::RustFS => {
let mut r_ = self.rustfs.as_ref().expect("err").clone(); if let Some(r_) = self.rustfs.as_ref() {
r_.secret_key = "REDACTED".to_string(); let mut r_clone = r_.clone();
r = Some(r_); r_clone.secret_key = "REDACTED".to_string();
r = Some(r_clone);
}
} }
TierType::MinIO => { TierType::MinIO => {
let mut compatible_backend_ = self.minio.as_ref().expect("err").clone(); if let Some(compatible_backend_) = self.minio.as_ref() {
compatible_backend_.secret_key = "REDACTED".to_string(); let mut compatible_backend_clone = compatible_backend_.clone();
compatible_backend = Some(compatible_backend_); compatible_backend_clone.secret_key = "REDACTED".to_string();
compatible_backend = Some(compatible_backend_clone);
}
} }
TierType::Aliyun => { TierType::Aliyun => {
let mut aliyun_ = self.aliyun.as_ref().expect("err").clone(); if let Some(aliyun_) = self.aliyun.as_ref() {
aliyun_.secret_key = "REDACTED".to_string(); let mut aliyun_clone = aliyun_.clone();
aliyun = Some(aliyun_); aliyun_clone.secret_key = "REDACTED".to_string();
aliyun = Some(aliyun_clone);
}
} }
TierType::Tencent => { TierType::Tencent => {
let mut tencent_ = self.tencent.as_ref().expect("err").clone(); if let Some(tencent_) = self.tencent.as_ref() {
tencent_.secret_key = "REDACTED".to_string(); let mut tencent_clone = tencent_.clone();
tencent = Some(tencent_); tencent_clone.secret_key = "REDACTED".to_string();
tencent = Some(tencent_clone);
}
} }
TierType::Huaweicloud => { TierType::Huaweicloud => {
let mut huaweicloud_ = self.huaweicloud.as_ref().expect("err").clone(); if let Some(huaweicloud_) = self.huaweicloud.as_ref() {
huaweicloud_.secret_key = "REDACTED".to_string(); let mut huaweicloud_clone = huaweicloud_.clone();
huaweicloud = Some(huaweicloud_); huaweicloud_clone.secret_key = "REDACTED".to_string();
huaweicloud = Some(huaweicloud_clone);
}
} }
TierType::Azure => { TierType::Azure => {
let mut azure_ = self.azure.as_ref().expect("err").clone(); if let Some(azure_) = self.azure.as_ref() {
azure_.secret_key = "REDACTED".to_string(); let mut azure_clone = azure_.clone();
azure = Some(azure_); azure_clone.secret_key = "REDACTED".to_string();
azure = Some(azure_clone);
}
} }
TierType::GCS => { TierType::GCS => {
let mut gcs_ = self.gcs.as_ref().expect("err").clone(); if let Some(gcs_) = self.gcs.as_ref() {
gcs_.creds = "REDACTED".to_string(); let mut gcs_clone = gcs_.clone();
gcs = Some(gcs_); gcs_clone.creds = "REDACTED".to_string();
gcs = Some(gcs_clone);
}
} }
TierType::R2 => { TierType::R2 => {
let mut r2_ = self.r2.as_ref().expect("err").clone(); if let Some(r2_) = self.r2.as_ref() {
r2_.secret_key = "REDACTED".to_string(); let mut r2_clone = r2_.clone();
r2 = Some(r2_); r2_clone.secret_key = "REDACTED".to_string();
r2 = Some(r2_clone);
}
} }
_ => (), _ => (),
} }
@@ -222,15 +240,15 @@ impl Clone for TierConfig {
impl TierConfig { impl TierConfig {
fn endpoint(&self) -> String { fn endpoint(&self) -> String {
match self.tier_type { match self.tier_type {
TierType::S3 => self.s3.as_ref().expect("err").endpoint.clone(), TierType::S3 => self.s3.as_ref().map(|s| s.endpoint.clone()).unwrap_or_default(),
TierType::RustFS => self.rustfs.as_ref().expect("err").endpoint.clone(), TierType::RustFS => self.rustfs.as_ref().map(|r| r.endpoint.clone()).unwrap_or_default(),
TierType::MinIO => self.minio.as_ref().expect("err").endpoint.clone(), TierType::MinIO => self.minio.as_ref().map(|m| m.endpoint.clone()).unwrap_or_default(),
TierType::Aliyun => self.aliyun.as_ref().expect("err").endpoint.clone(), TierType::Aliyun => self.aliyun.as_ref().map(|a| a.endpoint.clone()).unwrap_or_default(),
TierType::Tencent => self.tencent.as_ref().expect("err").endpoint.clone(), TierType::Tencent => self.tencent.as_ref().map(|t| t.endpoint.clone()).unwrap_or_default(),
TierType::Huaweicloud => self.huaweicloud.as_ref().expect("err").endpoint.clone(), TierType::Huaweicloud => self.huaweicloud.as_ref().map(|h| h.endpoint.clone()).unwrap_or_default(),
TierType::Azure => self.azure.as_ref().expect("err").endpoint.clone(), TierType::Azure => self.azure.as_ref().map(|a| a.endpoint.clone()).unwrap_or_default(),
TierType::GCS => self.gcs.as_ref().expect("err").endpoint.clone(), TierType::GCS => self.gcs.as_ref().map(|g| g.endpoint.clone()).unwrap_or_default(),
TierType::R2 => self.r2.as_ref().expect("err").endpoint.clone(), TierType::R2 => self.r2.as_ref().map(|r| r.endpoint.clone()).unwrap_or_default(),
_ => { _ => {
info!("unexpected tier type {}", self.tier_type); info!("unexpected tier type {}", self.tier_type);
"".to_string() "".to_string()
@@ -240,15 +258,15 @@ impl TierConfig {
fn bucket(&self) -> String { fn bucket(&self) -> String {
match self.tier_type { match self.tier_type {
TierType::S3 => self.s3.as_ref().expect("err").bucket.clone(), TierType::S3 => self.s3.as_ref().map(|s| s.bucket.clone()).unwrap_or_default(),
TierType::RustFS => self.rustfs.as_ref().expect("err").bucket.clone(), TierType::RustFS => self.rustfs.as_ref().map(|r| r.bucket.clone()).unwrap_or_default(),
TierType::MinIO => self.minio.as_ref().expect("err").bucket.clone(), TierType::MinIO => self.minio.as_ref().map(|m| m.bucket.clone()).unwrap_or_default(),
TierType::Aliyun => self.aliyun.as_ref().expect("err").bucket.clone(), TierType::Aliyun => self.aliyun.as_ref().map(|a| a.bucket.clone()).unwrap_or_default(),
TierType::Tencent => self.tencent.as_ref().expect("err").bucket.clone(), TierType::Tencent => self.tencent.as_ref().map(|t| t.bucket.clone()).unwrap_or_default(),
TierType::Huaweicloud => self.huaweicloud.as_ref().expect("err").bucket.clone(), TierType::Huaweicloud => self.huaweicloud.as_ref().map(|h| h.bucket.clone()).unwrap_or_default(),
TierType::Azure => self.azure.as_ref().expect("err").bucket.clone(), TierType::Azure => self.azure.as_ref().map(|a| a.bucket.clone()).unwrap_or_default(),
TierType::GCS => self.gcs.as_ref().expect("err").bucket.clone(), TierType::GCS => self.gcs.as_ref().map(|g| g.bucket.clone()).unwrap_or_default(),
TierType::R2 => self.r2.as_ref().expect("err").bucket.clone(), TierType::R2 => self.r2.as_ref().map(|r| r.bucket.clone()).unwrap_or_default(),
_ => { _ => {
info!("unexpected tier type {}", self.tier_type); info!("unexpected tier type {}", self.tier_type);
"".to_string() "".to_string()
@@ -258,15 +276,15 @@ impl TierConfig {
fn prefix(&self) -> String { fn prefix(&self) -> String {
match self.tier_type { match self.tier_type {
TierType::S3 => self.s3.as_ref().expect("err").prefix.clone(), TierType::S3 => self.s3.as_ref().map(|s| s.prefix.clone()).unwrap_or_default(),
TierType::RustFS => self.rustfs.as_ref().expect("err").prefix.clone(), TierType::RustFS => self.rustfs.as_ref().map(|r| r.prefix.clone()).unwrap_or_default(),
TierType::MinIO => self.minio.as_ref().expect("err").prefix.clone(), TierType::MinIO => self.minio.as_ref().map(|m| m.prefix.clone()).unwrap_or_default(),
TierType::Aliyun => self.aliyun.as_ref().expect("err").prefix.clone(), TierType::Aliyun => self.aliyun.as_ref().map(|a| a.prefix.clone()).unwrap_or_default(),
TierType::Tencent => self.tencent.as_ref().expect("err").prefix.clone(), TierType::Tencent => self.tencent.as_ref().map(|t| t.prefix.clone()).unwrap_or_default(),
TierType::Huaweicloud => self.huaweicloud.as_ref().expect("err").prefix.clone(), TierType::Huaweicloud => self.huaweicloud.as_ref().map(|h| h.prefix.clone()).unwrap_or_default(),
TierType::Azure => self.azure.as_ref().expect("err").prefix.clone(), TierType::Azure => self.azure.as_ref().map(|a| a.prefix.clone()).unwrap_or_default(),
TierType::GCS => self.gcs.as_ref().expect("err").prefix.clone(), TierType::GCS => self.gcs.as_ref().map(|g| g.prefix.clone()).unwrap_or_default(),
TierType::R2 => self.r2.as_ref().expect("err").prefix.clone(), TierType::R2 => self.r2.as_ref().map(|r| r.prefix.clone()).unwrap_or_default(),
_ => { _ => {
info!("unexpected tier type {}", self.tier_type); info!("unexpected tier type {}", self.tier_type);
"".to_string() "".to_string()
@@ -276,15 +294,15 @@ impl TierConfig {
fn region(&self) -> String { fn region(&self) -> String {
match self.tier_type { match self.tier_type {
TierType::S3 => self.s3.as_ref().expect("err").region.clone(), TierType::S3 => self.s3.as_ref().map(|s| s.region.clone()).unwrap_or_default(),
TierType::RustFS => self.rustfs.as_ref().expect("err").region.clone(), TierType::RustFS => self.rustfs.as_ref().map(|r| r.region.clone()).unwrap_or_default(),
TierType::MinIO => self.minio.as_ref().expect("err").region.clone(), TierType::MinIO => self.minio.as_ref().map(|m| m.region.clone()).unwrap_or_default(),
TierType::Aliyun => self.aliyun.as_ref().expect("err").region.clone(), TierType::Aliyun => self.aliyun.as_ref().map(|a| a.region.clone()).unwrap_or_default(),
TierType::Tencent => self.tencent.as_ref().expect("err").region.clone(), TierType::Tencent => self.tencent.as_ref().map(|t| t.region.clone()).unwrap_or_default(),
TierType::Huaweicloud => self.huaweicloud.as_ref().expect("err").region.clone(), TierType::Huaweicloud => self.huaweicloud.as_ref().map(|h| h.region.clone()).unwrap_or_default(),
TierType::Azure => self.azure.as_ref().expect("err").region.clone(), TierType::Azure => self.azure.as_ref().map(|a| a.region.clone()).unwrap_or_default(),
TierType::GCS => self.gcs.as_ref().expect("err").region.clone(), TierType::GCS => self.gcs.as_ref().map(|g| g.region.clone()).unwrap_or_default(),
TierType::R2 => self.r2.as_ref().expect("err").region.clone(), TierType::R2 => self.r2.as_ref().map(|r| r.region.clone()).unwrap_or_default(),
_ => { _ => {
info!("unexpected tier type {}", self.tier_type); info!("unexpected tier type {}", self.tier_type);
"".to_string() "".to_string()
+130 -52
View File
@@ -1,4 +1,3 @@
#![allow(unused_imports)]
// Copyright 2024 RustFS Team // Copyright 2024 RustFS Team
// //
// Licensed under the Apache License, Version 2.0 (the "License"); // Licensed under the Apache License, Version 2.0 (the "License");
@@ -12,6 +11,7 @@
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and // See the License for the specific language governing permissions and
// limitations under the License. // limitations under the License.
#![allow(unused_imports)]
#![allow(unused_variables)] #![allow(unused_variables)]
#![allow(unused_mut)] #![allow(unused_mut)]
#![allow(unused_assignments)] #![allow(unused_assignments)]
@@ -27,7 +27,7 @@ use crate::error::is_err_bucket_not_found;
use crate::tier::{ use crate::tier::{
tier::ERR_TIER_TYPE_UNSUPPORTED, tier::ERR_TIER_TYPE_UNSUPPORTED,
tier_config::{TierConfig, TierType}, tier_config::{TierConfig, TierType},
tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_PERM_ERR}, tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_NOT_FOUND, ERR_TIER_PERM_ERR},
warm_backend_aliyun::WarmBackendAliyun, warm_backend_aliyun::WarmBackendAliyun,
warm_backend_azure::WarmBackendAzure, warm_backend_azure::WarmBackendAzure,
warm_backend_gcs::WarmBackendGCS, warm_backend_gcs::WarmBackendGCS,
@@ -155,7 +155,7 @@ pub fn build_transition_put_options(storage_class: String, mut metadata: HashMap
} }
pub async fn check_warm_backend(w: Option<&WarmBackendImpl>) -> Result<(), AdminError> { pub async fn check_warm_backend(w: Option<&WarmBackendImpl>) -> Result<(), AdminError> {
let w = w.expect("err"); let w = w.ok_or_else(|| ERR_TIER_NOT_FOUND.clone())?;
let remote_version_id = w let remote_version_id = w
.put(PROBE_OBJECT, ReaderImpl::Body(Bytes::from("RustFS".as_bytes().to_vec())), 5) .put(PROBE_OBJECT, ReaderImpl::Body(Bytes::from("RustFS".as_bytes().to_vec())), 5)
.await; .await;
@@ -176,9 +176,11 @@ pub async fn check_warm_backend(w: Option<&WarmBackendImpl>) -> Result<(), Admin
return Err(ERR_TIER_PERM_ERR.clone()); return Err(ERR_TIER_PERM_ERR.clone());
//} //}
} }
if let Err(err) = w.remove(PROBE_OBJECT, &remote_version_id.expect("err")).await { if let Ok(version_id) = remote_version_id {
return Err(ERR_TIER_PERM_ERR.clone()); if let Err(err) = w.remove(PROBE_OBJECT, &version_id).await {
}; return Err(ERR_TIER_PERM_ERR.clone());
};
}
Ok(()) Ok(())
} }
@@ -186,119 +188,195 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
let mut d: Option<WarmBackendImpl> = None; let mut d: Option<WarmBackendImpl> = None;
match tier.tier_type { match tier.tier_type {
TierType::S3 => { TierType::S3 => {
let dd = WarmBackendS3::new(tier.s3.as_ref().expect("err"), &tier.name).await; if let Some(s3_config) = tier.s3.as_ref() {
if let Err(err) = dd { let dd = WarmBackendS3::new(s3_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create S3 backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "S3 tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::RustFS => { TierType::RustFS => {
let dd = WarmBackendRustFS::new(tier.rustfs.as_ref().expect("err"), &tier.name).await; if let Some(rustfs_config) = tier.rustfs.as_ref() {
if let Err(err) = dd { let dd = WarmBackendRustFS::new(rustfs_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create RustFS backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "RustFS tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::MinIO => { TierType::MinIO => {
let dd = WarmBackendMinIO::new(tier.minio.as_ref().expect("err"), &tier.name).await; if let Some(minio_config) = tier.minio.as_ref() {
if let Err(err) = dd { let dd = WarmBackendMinIO::new(minio_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create MinIO backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "MinIO tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::Aliyun => { TierType::Aliyun => {
let dd = WarmBackendAliyun::new(tier.aliyun.as_ref().expect("err"), &tier.name).await; if let Some(aliyun_config) = tier.aliyun.as_ref() {
if let Err(err) = dd { let dd = WarmBackendAliyun::new(aliyun_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create Aliyun backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "Aliyun tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::Tencent => { TierType::Tencent => {
let dd = WarmBackendTencent::new(tier.tencent.as_ref().expect("err"), &tier.name).await; if let Some(tencent_config) = tier.tencent.as_ref() {
if let Err(err) = dd { let dd = WarmBackendTencent::new(tencent_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create Tencent backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "Tencent tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::Huaweicloud => { TierType::Huaweicloud => {
let dd = WarmBackendHuaweicloud::new(tier.huaweicloud.as_ref().expect("err"), &tier.name).await; if let Some(huaweicloud_config) = tier.huaweicloud.as_ref() {
if let Err(err) = dd { let dd = WarmBackendHuaweicloud::new(huaweicloud_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create Huaweicloud backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "Huaweicloud tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::Azure => { TierType::Azure => {
let dd = WarmBackendAzure::new(tier.azure.as_ref().expect("err"), &tier.name).await; if let Some(azure_config) = tier.azure.as_ref() {
if let Err(err) = dd { let dd = WarmBackendAzure::new(azure_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create Azure backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "Azure tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::GCS => { TierType::GCS => {
let dd = WarmBackendGCS::new(tier.gcs.as_ref().expect("err"), &tier.name).await; if let Some(gcs_config) = tier.gcs.as_ref() {
if let Err(err) = dd { let dd = WarmBackendGCS::new(gcs_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create GCS backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "GCS tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
TierType::R2 => { TierType::R2 => {
let dd = WarmBackendR2::new(tier.r2.as_ref().expect("err"), &tier.name).await; if let Some(r2_config) = tier.r2.as_ref() {
if let Err(err) = dd { let dd = WarmBackendR2::new(r2_config, &tier.name).await;
warn!("{}", err); if let Err(err) = dd {
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
status_code: StatusCode::BAD_REQUEST,
});
}
d = Some(Box::new(dd.expect("Failed to create R2 backend")));
} else {
return Err(AdminError { return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(), code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()), message: "R2 tier configuration not found".to_string(),
status_code: StatusCode::BAD_REQUEST, status_code: StatusCode::BAD_REQUEST,
}); });
} }
d = Some(Box::new(dd.expect("err")));
} }
_ => { _ => {
return Err(ERR_TIER_TYPE_UNSUPPORTED.clone()); return Err(ERR_TIER_TYPE_UNSUPPORTED.clone());
} }
} }
Ok(d.expect("err")) d.ok_or_else(|| AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: "Tier backend not initialized".to_string(),
status_code: StatusCode::BAD_REQUEST,
})
} }
#[cfg(test)] #[cfg(test)]
@@ -76,12 +76,10 @@ impl WarmBackendAliyun {
}; };
let scheme = u.scheme(); let scheme = u.scheme();
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let client = TransitionClient::new( let host = u
&format!("{}:{}", u.host_str().expect("err"), u.port().unwrap_or(default_port)), .host_str()
opts, .ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
"aliyun", let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "aliyun").await?;
)
.await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));
@@ -76,12 +76,10 @@ impl WarmBackendAzure {
}; };
let scheme = u.scheme(); let scheme = u.scheme();
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let client = TransitionClient::new( let host = u
&format!("{}:{}", u.host_str().expect("err"), u.port().unwrap_or(default_port)), .host_str()
opts, .ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
"azure", let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "azure").await?;
)
.await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));
@@ -76,12 +76,11 @@ impl WarmBackendHuaweicloud {
}; };
let scheme = u.scheme(); let scheme = u.scheme();
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let client = TransitionClient::new( let host = u
&format!("{}:{}", u.host_str().expect("err"), u.port().unwrap_or(default_port)), .host_str()
opts, .ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
"huaweicloud", let client =
) TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "huaweicloud").await?;
.await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));
@@ -75,12 +75,10 @@ impl WarmBackendMinIO {
}; };
let scheme = u.scheme(); let scheme = u.scheme();
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let client = TransitionClient::new( let host = u
&format!("{}:{}", u.host_str().expect("err"), u.port().unwrap_or(default_port)), .host_str()
opts, .ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
"minio", let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "minio").await?;
)
.await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));
+4 -6
View File
@@ -75,12 +75,10 @@ impl WarmBackendR2 {
}; };
let scheme = u.scheme(); let scheme = u.scheme();
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let client = TransitionClient::new( let host = u
&format!("{}:{}", u.host_str().expect("err"), u.port().unwrap_or(default_port)), .host_str()
opts, .ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
"r2", let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "r2").await?;
)
.await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));
+8 -3
View File
@@ -95,7 +95,10 @@ impl WarmBackendS3 {
region: conf.region.clone(), region: conf.region.clone(),
..Default::default() ..Default::default()
}; };
let client = TransitionClient::new(&u.host().expect("err").to_string(), opts, "s3").await?; let host = u
.host()
.ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
let client = TransitionClient::new(&host.to_string(), opts, "s3").await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));
@@ -164,8 +167,10 @@ impl WarmBackend for WarmBackendS3 {
ropts.version_id = rv.to_string(); ropts.version_id = rv.to_string();
} }
let client = self.client.clone(); let client = self.client.clone();
let err = client.remove_object(&self.bucket, &self.get_dest(object), ropts).await; match client.remove_object(&self.bucket, &self.get_dest(object), ropts).await {
Err(std::io::Error::other(err.expect("err"))) None => Ok(()),
Some(err) => Err(std::io::Error::other(err)),
}
} }
async fn in_use(&self) -> Result<bool, std::io::Error> { async fn in_use(&self) -> Result<bool, std::io::Error> {
@@ -190,6 +190,6 @@ impl WarmBackend for WarmBackendS3 {
return Err(std::io::Error::other("list_objects_v2 error")); return Err(std::io::Error::other("list_objects_v2 error"));
}; };
Ok(res.common_prefixes.unwrap().len() > 0 || res.contents.unwrap().len() > 0) Ok(res.common_prefixes.unwrap_or_default().len() > 0 || res.contents.unwrap_or_default().len() > 0)
} }
} }
@@ -76,12 +76,10 @@ impl WarmBackendTencent {
}; };
let scheme = u.scheme(); let scheme = u.scheme();
let default_port = if scheme == "https" { 443 } else { 80 }; let default_port = if scheme == "https" { 443 } else { 80 };
let client = TransitionClient::new( let host = u
&format!("{}:{}", u.host_str().expect("err"), u.port().unwrap_or(default_port)), .host_str()
opts, .ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
"tencent", let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "tencent").await?;
)
.await?;
let client = Arc::new(client); let client = Arc::new(client);
let core = TransitionCore(Arc::clone(&client)); let core = TransitionCore(Arc::clone(&client));