From 8b57076194290145130153ff11be786bac24d651 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 12 Aug 2026 19:03:31 +0800 Subject: [PATCH] chore(ecstore): remove dead MinIO-port client modules (~940 lines) (#5982) Delete five zero-caller client modules (api_bucket_policy, api_get_object_acl, api_get_object_attributes, api_get_object_file, api_restore), the orphaned TransitionCore get/put_bucket_policy wrappers, and the two constants only they consumed. Also removes two unwrap() panics on remote-controlled data in api_get_object_acl.rs (non-UTF-8 body, missing Owner ID). Tier warm-backend APIs untouched. Ref: rustfs/backlog#1822 (T1) --- .../ecstore/src/client/api_bucket_policy.rs | 171 ----------- .../ecstore/src/client/api_get_object_acl.rs | 199 ------------- .../src/client/api_get_object_attributes.rs | 266 ------------------ .../ecstore/src/client/api_get_object_file.rs | 159 ----------- crates/ecstore/src/client/api_restore.rs | 134 --------- crates/ecstore/src/client/constants.rs | 3 - crates/ecstore/src/client/mod.rs | 5 - crates/ecstore/src/client/transition_api.rs | 10 - 8 files changed, 947 deletions(-) delete mode 100644 crates/ecstore/src/client/api_bucket_policy.rs delete mode 100644 crates/ecstore/src/client/api_get_object_acl.rs delete mode 100644 crates/ecstore/src/client/api_get_object_attributes.rs delete mode 100644 crates/ecstore/src/client/api_get_object_file.rs delete mode 100644 crates/ecstore/src/client/api_restore.rs diff --git a/crates/ecstore/src/client/api_bucket_policy.rs b/crates/ecstore/src/client/api_bucket_policy.rs deleted file mode 100644 index c01039508..000000000 --- a/crates/ecstore/src/client/api_bucket_policy.rs +++ /dev/null @@ -1,171 +0,0 @@ -// 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. -#![allow(unused_imports)] -#![allow(unused_variables)] -#![allow(unused_mut)] -#![allow(unused_assignments)] -#![allow(unused_must_use)] -#![allow(clippy::all)] - -use http::{HeaderMap, StatusCode}; -use http_body_util::BodyExt; -use hyper::body::Body; -use hyper::body::Bytes; -use std::collections::HashMap; - -use crate::client::{ - api_error_response::http_resp_to_error_response, - transition_api::{ReaderImpl, RequestMetadata, TransitionClient}, -}; -use rustfs_utils::hash::EMPTY_STRING_SHA256_HASH; - -impl TransitionClient { - pub async fn set_bucket_policy(&self, bucket_name: &str, policy: &str) -> Result<(), std::io::Error> { - if policy == "" { - return self.remove_bucket_policy(bucket_name).await; - } - - self.put_bucket_policy(bucket_name, policy).await - } - - pub async fn put_bucket_policy(&self, bucket_name: &str, policy: &str) -> Result<(), std::io::Error> { - let mut url_values = HashMap::new(); - url_values.insert("policy".to_string(), "".to_string()); - - let mut req_metadata = RequestMetadata { - bucket_name: bucket_name.to_string(), - query_values: url_values, - content_body: ReaderImpl::Body(Bytes::from(policy.as_bytes().to_vec())), - content_length: policy.len() as i64, - object_name: "".to_string(), - custom_header: HeaderMap::new(), - content_md5_base64: "".to_string(), - content_sha256_hex: "".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(), - }; - - let resp = self.execute_method(http::Method::PUT, &mut req_metadata).await?; - //defer closeResponse(resp) - - let resp_status = resp.status(); - let h = resp.headers().clone(); - - //if resp != nil { - if resp_status != StatusCode::NO_CONTENT && resp.status() != StatusCode::OK { - return Err(std::io::Error::other(http_resp_to_error_response( - resp_status, - &h, - vec![], - bucket_name, - "", - ))); - } - //} - Ok(()) - } - - pub async fn remove_bucket_policy(&self, bucket_name: &str) -> Result<(), std::io::Error> { - let mut url_values = HashMap::new(); - 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(), - }, - ) - .await?; - //defer closeResponse(resp) - - let resp_status = resp.status(); - let h = resp.headers().clone(); - - if resp_status != StatusCode::NO_CONTENT { - return Err(std::io::Error::other(http_resp_to_error_response( - resp_status, - &h, - vec![], - bucket_name, - "", - ))); - } - - Ok(()) - } - - pub async fn get_bucket_policy(&self, bucket_name: &str) -> Result { - let bucket_policy = self.get_bucket_policy_inner(bucket_name).await?; - Ok(bucket_policy) - } - - pub async fn get_bucket_policy_inner(&self, bucket_name: &str) -> Result { - let mut url_values = HashMap::new(); - 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(), - }, - ) - .await?; - - let mut body_vec = Vec::new(); - let mut body = resp.into_body(); - while let Some(frame) = body.frame().await { - let frame = frame.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e.to_string()))?; - if let Some(data) = frame.data_ref() { - body_vec.extend_from_slice(data); - } - } - let policy = String::from_utf8_lossy(&body_vec).to_string(); - Ok(policy) - } -} diff --git a/crates/ecstore/src/client/api_get_object_acl.rs b/crates/ecstore/src/client/api_get_object_acl.rs deleted file mode 100644 index bec3d8e2c..000000000 --- a/crates/ecstore/src/client/api_get_object_acl.rs +++ /dev/null @@ -1,199 +0,0 @@ -// 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. -#![allow(unused_imports)] -#![allow(unused_variables)] -#![allow(unused_mut)] -#![allow(unused_assignments)] -#![allow(unused_must_use)] -#![allow(clippy::all)] - -use crate::client::{ - api_error_response::http_resp_to_error_response, - api_get_options::GetObjectOptions, - transition_api::{ObjectInfo, ReaderImpl, RequestMetadata, TransitionClient}, -}; -use bytes::Bytes; -use http::{HeaderMap, HeaderValue}; -use http_body_util::BodyExt; -use rustfs_config::MAX_S3_CLIENT_RESPONSE_SIZE; -use rustfs_utils::EMPTY_STRING_SHA256_HASH; -use s3s::dto::Owner; -use std::collections::HashMap; - -#[derive(Clone, Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct Grantee { - pub id: String, - pub display_name: String, - pub uri: String, -} - -#[derive(Clone, Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct Grant { - pub grantee: Grantee, - pub permission: String, -} - -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct AccessControlList { - pub grant: Vec, - pub permission: String, -} - -#[derive(Debug, Default, serde::Deserialize)] -pub struct AccessControlPolicy { - #[serde(skip)] - owner: Owner, - pub access_control_list: AccessControlList, -} - -impl TransitionClient { - pub async fn get_object_acl(&self, bucket_name: &str, object_name: &str) -> Result { - 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(), - }, - ) - .await?; - - let resp_status = resp.status(); - let h = resp.headers().clone(); - - let mut body_vec = Vec::new(); - let mut body = resp.into_body(); - while let Some(frame) = body.frame().await { - let frame = frame.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e.to_string()))?; - if let Some(data) = frame.data_ref() { - body_vec.extend_from_slice(data); - } - } - - if resp_status != http::StatusCode::OK { - return Err(std::io::Error::other(http_resp_to_error_response( - resp_status, - &h, - body_vec, - bucket_name, - object_name, - ))); - } - - let mut res = match quick_xml::de::from_str::(&String::from_utf8(body_vec).unwrap()) { - Ok(result) => result, - Err(err) => { - return Err(std::io::Error::other(err.to_string())); - } - }; - - let mut obj_info = self - .stat_object(bucket_name, object_name, &GetObjectOptions::default()) - .await?; - - obj_info.owner.display_name = res.owner.display_name.clone(); - obj_info.owner.id = res.owner.id.clone(); - - //obj_info.grant.extend(res.access_control_list.grant); - - let canned_acl = get_canned_acl(&res); - if canned_acl != "" { - obj_info - .metadata - .insert("X-Amz-Acl", HeaderValue::from_str(&canned_acl).unwrap()); - return Ok(obj_info); - } - - let grant_acl = get_amz_grant_acl(&res); - /*for (k, v) in grant_acl { - obj_info.metadata.insert(HeaderName::from_bytes(k.as_bytes()).unwrap(), HeaderValue::from_str(&v.to_string()).unwrap()); - }*/ - - Ok(obj_info) - } -} - -fn get_canned_acl(ac_policy: &AccessControlPolicy) -> String { - let grants = ac_policy.access_control_list.grant.clone(); - - if grants.len() == 1 { - if grants[0].grantee.uri == "" && grants[0].permission == "FULL_CONTROL" { - return "private".to_string(); - } - } else if grants.len() == 2 { - for g in grants { - if g.grantee.uri == "http://acs.amazonaws.com/groups/global/AuthenticatedUsers" && &g.permission == "READ" { - return "authenticated-read".to_string(); - } - if g.grantee.uri == "http://acs.amazonaws.com/groups/global/AllUsers" && &g.permission == "READ" { - return "public-read".to_string(); - } - if g.permission == "READ" && g.grantee.id == ac_policy.owner.id.clone().unwrap() { - return "bucket-owner-read".to_string(); - } - } - } else if grants.len() == 3 { - for g in grants { - if g.grantee.uri == "http://acs.amazonaws.com/groups/global/AllUsers" && g.permission == "WRITE" { - return "public-read-write".to_string(); - } - } - } - "".to_string() -} - -pub fn get_amz_grant_acl(ac_policy: &AccessControlPolicy) -> HashMap> { - let grants = ac_policy.access_control_list.grant.clone(); - let mut res = HashMap::>::new(); - - for g in grants { - let mut id = "id=".to_string(); - id.push_str(&g.grantee.id); - let permission: &str = &g.permission; - match permission { - "READ" => { - res.entry("X-Amz-Grant-Read".to_string()).or_insert(vec![]).push(id); - } - "WRITE" => { - res.entry("X-Amz-Grant-Write".to_string()).or_insert(vec![]).push(id); - } - "READ_ACP" => { - res.entry("X-Amz-Grant-Read-Acp".to_string()).or_insert(vec![]).push(id); - } - "WRITE_ACP" => { - res.entry("X-Amz-Grant-Write-Acp".to_string()).or_insert(vec![]).push(id); - } - "FULL_CONTROL" => { - res.entry("X-Amz-Grant-Full-Control".to_string()).or_insert(vec![]).push(id); - } - _ => (), - } - } - res -} diff --git a/crates/ecstore/src/client/api_get_object_attributes.rs b/crates/ecstore/src/client/api_get_object_attributes.rs deleted file mode 100644 index eb6c55be6..000000000 --- a/crates/ecstore/src/client/api_get_object_attributes.rs +++ /dev/null @@ -1,266 +0,0 @@ -// 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. -#![allow(unused_imports)] -#![allow(unused_variables)] -#![allow(unused_mut)] -#![allow(unused_assignments)] -#![allow(unused_must_use)] -#![allow(clippy::all)] - -use http::{HeaderMap, HeaderValue}; -use std::collections::HashMap; -use time::OffsetDateTime; - -use crate::client::constants::{GET_OBJECT_ATTRIBUTES_MAX_PARTS, GET_OBJECT_ATTRIBUTES_TAGS, ISO8601_DATEFORMAT}; -use crate::client::{ - api_get_object_acl::AccessControlPolicy, - transition_api::{ReaderImpl, RequestMetadata, TransitionClient}, -}; -use http_body_util::BodyExt; -use hyper::body::Body; -use hyper::body::Bytes; -use hyper::body::Incoming; -use rustfs_config::MAX_S3_CLIENT_RESPONSE_SIZE; -use rustfs_utils::EMPTY_STRING_SHA256_HASH; -use s3s::header::{X_AMZ_MAX_PARTS, X_AMZ_OBJECT_ATTRIBUTES, X_AMZ_PART_NUMBER_MARKER, X_AMZ_VERSION_ID}; - -pub struct ObjectAttributesOptions { - pub max_parts: i64, - pub version_id: String, - pub part_number_marker: i64, - //server_side_encryption: encrypt::ServerSide, -} - -pub struct ObjectAttributes { - pub version_id: String, - pub last_modified: OffsetDateTime, - pub object_attributes_response: ObjectAttributesResponse, -} - -impl ObjectAttributes { - fn new() -> Self { - Self { - version_id: "".to_string(), - last_modified: OffsetDateTime::now_utc(), - object_attributes_response: ObjectAttributesResponse::new(), - } - } -} - -#[derive(Debug, Default, serde::Deserialize)] -pub struct Checksum { - checksum_crc32: String, - checksum_crc32c: String, - checksum_sha1: String, - checksum_sha256: String, -} - -impl Checksum { - fn new() -> Self { - Self { - checksum_crc32: "".to_string(), - checksum_crc32c: "".to_string(), - checksum_sha1: "".to_string(), - checksum_sha256: "".to_string(), - } - } -} - -#[derive(Debug, Default, serde::Deserialize)] -pub struct ObjectParts { - pub parts_count: i64, - pub part_number_marker: i64, - pub next_part_number_marker: i64, - pub max_parts: i64, - is_truncated: bool, - parts: Vec, -} - -impl ObjectParts { - fn new() -> Self { - Self { - parts_count: 0, - part_number_marker: 0, - next_part_number_marker: 0, - max_parts: 0, - is_truncated: false, - parts: Vec::new(), - } - } -} - -#[derive(Debug, Default, serde::Deserialize)] -pub struct ObjectAttributesResponse { - pub etag: String, - pub storage_class: String, - pub object_size: i64, - pub checksum: Checksum, - pub object_parts: ObjectParts, -} - -impl ObjectAttributesResponse { - fn new() -> Self { - Self { - etag: "".to_string(), - storage_class: "".to_string(), - object_size: 0, - checksum: Checksum::new(), - object_parts: ObjectParts::new(), - } - } -} - -#[derive(Debug, Default, serde::Deserialize)] -struct ObjectAttributePart { - checksum_crc32: String, - checksum_crc32c: String, - checksum_sha1: String, - checksum_sha256: String, - part_number: i64, - size: i64, -} - -impl ObjectAttributes { - pub async fn parse_response(&mut self, h: &HeaderMap, body_vec: Vec) -> Result<(), std::io::Error> { - let last_modified = h - .get("Last-Modified") - .ok_or_else(|| std::io::Error::other("missing Last-Modified header"))? - .to_str() - .map_err(|e| std::io::Error::other(format!("invalid Last-Modified header: {e}")))?; - let mod_time = OffsetDateTime::parse(last_modified, ISO8601_DATEFORMAT) - .map_err(|e| std::io::Error::other(format!("invalid Last-Modified date: {e}")))?; - self.last_modified = mod_time; - - let version_id = h - .get(X_AMZ_VERSION_ID) - .ok_or_else(|| std::io::Error::other("missing version ID header"))? - .to_str() - .map_err(|e| std::io::Error::other(format!("invalid version ID header: {e}")))?; - self.version_id = version_id.to_string(); - - let body_str = String::from_utf8(body_vec).map_err(|e| std::io::Error::other(format!("invalid UTF-8 body: {e}")))?; - let mut response = match quick_xml::de::from_str::(&body_str) { - Ok(result) => result, - Err(err) => { - return Err(std::io::Error::other(err.to_string())); - } - }; - self.object_attributes_response = response; - - Ok(()) - } -} - -impl TransitionClient { - pub async fn get_object_attributes( - &self, - bucket_name: &str, - object_name: &str, - opts: ObjectAttributesOptions, - ) -> Result { - let mut url_values = HashMap::new(); - url_values.insert("attributes".to_string(), "".to_string()); - if opts.version_id != "" { - url_values.insert("versionId".to_string(), opts.version_id); - } - - let mut headers = HeaderMap::new(); - headers.insert( - X_AMZ_OBJECT_ATTRIBUTES, - HeaderValue::from_str(GET_OBJECT_ATTRIBUTES_TAGS).expect("valid header value"), - ); - - if opts.part_number_marker > 0 { - headers.insert( - X_AMZ_PART_NUMBER_MARKER, - HeaderValue::from_str(&opts.part_number_marker.to_string()).expect("valid header value"), - ); - } - - if opts.max_parts > 0 { - headers.insert( - X_AMZ_MAX_PARTS, - HeaderValue::from_str(&opts.max_parts.to_string()).expect("valid header value"), - ); - } else { - headers.insert( - X_AMZ_MAX_PARTS, - HeaderValue::from_str(&GET_OBJECT_ATTRIBUTES_MAX_PARTS.to_string()).expect("valid header value"), - ); - } - - /*if opts.server_side_encryption.is_some() { - opts.server_side_encryption.Marshal(headers); - }*/ - - 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(), - }, - ) - .await?; - - let resp_status = resp.status(); - let h = resp.headers().clone(); - let has_etag = h.get("ETag").and_then(|v| v.to_str().ok()).unwrap_or(""); - if !has_etag.is_empty() { - return Err(std::io::Error::other( - "get_object_attributes is not supported by the current endpoint version", - )); - } - - let mut body_vec = Vec::new(); - let mut body = resp.into_body(); - while let Some(frame) = body.frame().await { - let frame = frame.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e.to_string()))?; - if let Some(data) = frame.data_ref() { - body_vec.extend_from_slice(data); - } - } - - if resp_status != http::StatusCode::OK { - let err_body = - String::from_utf8(body_vec).map_err(|e| std::io::Error::other(format!("invalid UTF-8 error body: {e}")))?; - let mut er = match quick_xml::de::from_str::(&err_body) { - Ok(result) => result, - Err(err) => { - return Err(std::io::Error::other(err.to_string())); - } - }; - - return Err(std::io::Error::other(er.access_control_list.permission)); - } - - let mut oa = ObjectAttributes::new(); - oa.parse_response(&h, body_vec).await?; - - Ok(oa) - } -} diff --git a/crates/ecstore/src/client/api_get_object_file.rs b/crates/ecstore/src/client/api_get_object_file.rs deleted file mode 100644 index 05694e274..000000000 --- a/crates/ecstore/src/client/api_get_object_file.rs +++ /dev/null @@ -1,159 +0,0 @@ -// 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. - -use std::io; -use std::path::{Path, PathBuf}; - -#[cfg(not(windows))] -use std::os::unix::fs::PermissionsExt; - -use tokio::fs::{self, OpenOptions}; -use tokio::io::{AsyncSeekExt, AsyncWriteExt, SeekFrom}; - -use crate::client::{ - api_error_response::err_invalid_argument, api_get_options::GetObjectOptions, transition_api::TransitionClient, -}; - -async fn prepare_download_target(file_path: &Path) -> io::Result<()> { - match fs::metadata(file_path).await { - Ok(metadata) if metadata.is_dir() => { - return Err(io::Error::other(err_invalid_argument("filename is a directory."))); - } - Ok(_) => {} - Err(err) if err.kind() == io::ErrorKind::NotFound => {} - Err(err) => return Err(err), - } - - if let Some(parent) = file_path.parent() - && !parent.as_os_str().is_empty() - { - fs::create_dir_all(parent).await?; - - #[cfg(not(windows))] - { - let mut permissions = fs::metadata(parent).await?.permissions(); - permissions.set_mode(0o700); - fs::set_permissions(parent, permissions).await?; - } - } - - Ok(()) -} - -fn build_part_path(file_path: &Path) -> PathBuf { - PathBuf::from(format!("{}.part.rustfs", file_path.display())) -} - -async fn open_download_part_file(file_part_path: &Path) -> io::Result { - let mut options = OpenOptions::new(); - options.create(true).truncate(false).read(true).write(true); - - #[cfg(not(windows))] - options.mode(0o600); - - options.open(file_part_path).await -} - -async fn cleanup_part_file(file_part_path: &Path) { - let _ = fs::remove_file(file_part_path).await; -} - -impl TransitionClient { - pub async fn fget_object( - &self, - bucket_name: &str, - object_name: &str, - file_path: &str, - mut opts: GetObjectOptions, - ) -> Result<(), io::Error> { - let file_path = Path::new(file_path); - prepare_download_target(file_path).await?; - - let file_part_path = build_part_path(file_path); - let mut file_part = open_download_part_file(&file_part_path).await?; - let existing_len = file_part.metadata().await?.len(); - if existing_len > 0 { - opts.set_range(existing_len as i64, 0)?; - file_part.seek(SeekFrom::Start(existing_len)).await?; - } - - let (_object_info, _headers, mut object_reader) = self.get_object_inner(bucket_name, object_name, &opts).await?; - if let Err(err) = tokio::io::copy(&mut object_reader, &mut file_part).await { - cleanup_part_file(&file_part_path).await; - return Err(err); - } - - if let Err(err) = file_part.flush().await { - cleanup_part_file(&file_part_path).await; - return Err(err); - } - drop(file_part); - - if let Err(err) = fs::rename(&file_part_path, file_path).await { - cleanup_part_file(&file_part_path).await; - return Err(err); - } - - Ok(()) - } -} - -#[cfg(test)] -mod tests { - use super::*; - use tempfile::tempdir; - - #[tokio::test] - async fn prepare_download_target_allows_missing_file_and_creates_parent_dirs() { - let dir = tempdir().expect("temp dir"); - let target = dir.path().join("nested").join("object.bin"); - - prepare_download_target(&target) - .await - .expect("missing target should be accepted"); - - assert!(target.parent().expect("parent").exists(), "parent directory should be created"); - assert!( - fs::metadata(&target).await.is_err(), - "preparing the target should not create the final file eagerly" - ); - } - - #[tokio::test] - async fn prepare_download_target_rejects_directory_paths() { - let dir = tempdir().expect("temp dir"); - let target_dir = dir.path().join("download-dir"); - fs::create_dir_all(&target_dir).await.expect("target dir"); - - let err = prepare_download_target(&target_dir) - .await - .expect_err("directory targets must be rejected"); - - assert!(err.to_string().contains("directory"), "unexpected error for directory target: {err}"); - } - - #[tokio::test] - async fn open_download_part_file_creates_part_file() { - let dir = tempdir().expect("temp dir"); - let target = dir.path().join("object.bin"); - let part_path = build_part_path(&target); - - let file = open_download_part_file(&part_path) - .await - .expect("part file should be created"); - drop(file); - - assert!(part_path.exists(), "part file should exist after creation"); - } -} diff --git a/crates/ecstore/src/client/api_restore.rs b/crates/ecstore/src/client/api_restore.rs deleted file mode 100644 index 04e76bf21..000000000 --- a/crates/ecstore/src/client/api_restore.rs +++ /dev/null @@ -1,134 +0,0 @@ -// 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. -#![allow(unused_imports)] -#![allow(unused_variables)] -#![allow(unused_mut)] -#![allow(unused_assignments)] -#![allow(unused_must_use)] -#![allow(clippy::all)] - -use crate::client::{ - api_error_response::{err_invalid_argument, http_resp_to_error_response}, - api_get_object_acl::AccessControlList, - api_get_options::GetObjectOptions, - transition_api::{ObjectInfo, ReadCloser, ReaderImpl, RequestMetadata, TransitionClient, to_object_info}, -}; -use http::HeaderMap; -use http_body_util::BodyExt; -use hyper::body::Body; -use hyper::body::Bytes; -use s3s::dto::RestoreRequest; -use std::collections::HashMap; -use std::io::Cursor; -use tokio::io::BufReader; - -const TIER_STANDARD: &str = "Standard"; -const TIER_BULK: &str = "Bulk"; -const TIER_EXPEDITED: &str = "Expedited"; - -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct Encryption { - pub encryption_type: String, - pub kms_context: String, - pub kms_key_id: String, -} - -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct MetadataEntry { - pub name: String, - pub value: String, -} - -#[derive(Debug, Default, serde::Serialize)] -pub struct S3 { - pub access_control_list: AccessControlList, - pub bucket_name: String, - pub prefix: String, - pub canned_acl: String, - pub encryption: Encryption, - pub storage_class: String, - //tagging: Tags, - pub user_metadata: MetadataEntry, -} - -impl TransitionClient { - pub async fn restore_object( - &self, - bucket_name: &str, - object_name: &str, - version_id: &str, - restore_req: &RestoreRequest, - ) -> Result<(), std::io::Error> { - /*let restore_request = match quick_xml::se::to_string(restore_req) { - Ok(buf) => buf, - Err(e) => { - return Err(std::io::Error::other(e)); - } - };*/ - let restore_request = "".to_string(); - let restore_request_bytes = restore_request.as_bytes().to_vec(); - - let mut url_values = HashMap::new(); - url_values.insert("restore".to_string(), "".to_string()); - if version_id != "" { - url_values.insert("versionId".to_string(), version_id.to_string()); - } - - 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(), - }, - ) - .await?; - - let resp_status = resp.status(); - let h = resp.headers().clone(); - - let mut body_vec = Vec::new(); - let mut body = resp.into_body(); - while let Some(frame) = body.frame().await { - let frame = frame.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e.to_string()))?; - if let Some(data) = frame.data_ref() { - body_vec.extend_from_slice(data); - } - } - if resp_status != http::StatusCode::ACCEPTED && resp_status != http::StatusCode::OK { - return Err(std::io::Error::other(http_resp_to_error_response( - resp_status, - &h, - body_vec, - bucket_name, - "", - ))); - } - Ok(()) - } -} diff --git a/crates/ecstore/src/client/constants.rs b/crates/ecstore/src/client/constants.rs index 5324e159c..5c5a02482 100644 --- a/crates/ecstore/src/client/constants.rs +++ b/crates/ecstore/src/client/constants.rs @@ -37,6 +37,3 @@ pub const TOTAL_WORKERS: i64 = 4; pub const SIGN_V4_ALGORITHM: &str = "AWS4-HMAC-SHA256"; pub const ISO8601_DATEFORMAT: &[FormatItem<'_>] = format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond]Z"); - -pub const GET_OBJECT_ATTRIBUTES_TAGS: &str = "ETag,Checksum,StorageClass,ObjectSize,ObjectParts"; -pub const GET_OBJECT_ATTRIBUTES_MAX_PARTS: i64 = 1000; diff --git a/crates/ecstore/src/client/mod.rs b/crates/ecstore/src/client/mod.rs index eebbe42cc..39902cb8a 100644 --- a/crates/ecstore/src/client/mod.rs +++ b/crates/ecstore/src/client/mod.rs @@ -16,12 +16,8 @@ #![allow(dead_code)] pub mod admin_handler_utils; -pub mod api_bucket_policy; pub mod api_error_response; pub mod api_get_object; -pub mod api_get_object_acl; -pub mod api_get_object_attributes; -pub mod api_get_object_file; pub mod api_get_options; pub mod api_list; pub mod api_put_object; @@ -29,7 +25,6 @@ pub mod api_put_object_common; pub mod api_put_object_multipart; pub mod api_put_object_streaming; pub mod api_remove; -pub mod api_restore; pub mod api_s3_datatypes; pub mod api_stat; pub mod bucket_cache; diff --git a/crates/ecstore/src/client/transition_api.rs b/crates/ecstore/src/client/transition_api.rs index 53153c421..38e65b7ab 100644 --- a/crates/ecstore/src/client/transition_api.rs +++ b/crates/ecstore/src/client/transition_api.rs @@ -1006,16 +1006,6 @@ impl TransitionCore { client.abort_multipart_upload(bucket_name, object, upload_id).await } - pub async fn get_bucket_policy(&self, bucket_name: &str) -> Result { - let client = self.0.clone(); - client.get_bucket_policy(bucket_name).await - } - - pub async fn put_bucket_policy(&self, bucket_name: &str, bucket_policy: &str) -> Result<(), std::io::Error> { - let client = self.0.clone(); - client.put_bucket_policy(bucket_name, bucket_policy).await - } - pub async fn get_object( &self, bucket_name: &str,