// 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, Method, StatusCode}; use http_body_util::BodyExt; use hyper::body::Body; use hyper::body::Bytes; use rustfs_utils::HashAlgorithm; use s3s::S3ErrorCode; use s3s::dto::ReplicationStatus; use s3s::header::{X_AMZ_BYPASS_GOVERNANCE_RETENTION, X_AMZ_DELETE_MARKER, X_AMZ_VERSION_ID}; use serde::Deserialize; use std::fmt::Display; use std::{ collections::{HashMap, HashSet}, sync::Arc, }; use time::OffsetDateTime; use tokio::sync::mpsc::{self, Receiver, Sender}; use crate::client::utils::base64_encode; use crate::client::{ api_error_response::{ErrorResponse, http_resp_to_error_response, to_error_response}, api_s3_datatypes::{DeleteMultiObjects, DeleteObject}, transition_api::{ReaderImpl, RequestMetadata, TransitionClient}, }; use crate::{ disk::DiskAPI, store_api::{GetObjectReader, ObjectInfo, StorageAPI}, }; use rustfs_utils::hash::EMPTY_STRING_SHA256_HASH; pub struct RemoveBucketOptions { _forced_delete: bool, } const DELETE_RESPONSE_PREVIEW_LEN: usize = 1024; #[derive(Debug)] #[allow(dead_code)] pub struct AdvancedRemoveOptions { pub replication_delete_marker: bool, pub replication_status: ReplicationStatus, pub replication_mtime: Option, pub replication_request: bool, pub replication_validity_check: bool, } impl Default for AdvancedRemoveOptions { fn default() -> Self { Self { replication_delete_marker: false, replication_status: ReplicationStatus::from_static(ReplicationStatus::PENDING), replication_mtime: None, replication_request: false, replication_validity_check: false, } } } #[derive(Debug, Default)] pub struct RemoveObjectOptions { pub force_delete: bool, pub governance_bypass: bool, pub version_id: String, pub internal: AdvancedRemoveOptions, } impl TransitionClient { pub async fn remove_bucket_with_options(&self, bucket_name: &str, opts: &RemoveBucketOptions) -> Result<(), std::io::Error> { let headers = HeaderMap::new(); let resp = self .execute_method( Method::DELETE, &mut RequestMetadata { bucket_name: bucket_name.to_string(), content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), custom_header: headers, object_name: "".to_string(), query_values: Default::default(), content_body: ReaderImpl::Body(Bytes::new()), content_length: 0, content_md5_base64: "".to_string(), stream_sha256: false, trailer: HeaderMap::new(), pre_sign_url: Default::default(), add_crc: Default::default(), extra_pre_sign_header: Default::default(), bucket_location: Default::default(), expires: Default::default(), }, ) .await?; { if let Ok(mut bucket_loc_cache) = self.bucket_loc_cache.lock() { bucket_loc_cache.delete(bucket_name); } } Ok(()) } pub async fn remove_bucket(&self, bucket_name: &str) -> Result<(), std::io::Error> { let resp = self .execute_method( http::Method::DELETE, &mut RequestMetadata { bucket_name: bucket_name.to_string(), content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), custom_header: Default::default(), object_name: "".to_string(), query_values: Default::default(), content_body: ReaderImpl::Body(Bytes::new()), content_length: 0, content_md5_base64: "".to_string(), stream_sha256: false, trailer: HeaderMap::new(), pre_sign_url: Default::default(), add_crc: Default::default(), extra_pre_sign_header: Default::default(), bucket_location: Default::default(), expires: Default::default(), }, ) .await?; { if let Ok(mut bucket_loc_cache) = self.bucket_loc_cache.lock() { bucket_loc_cache.delete(bucket_name); } } Ok(()) } pub async fn remove_object(&self, bucket_name: &str, object_name: &str, opts: RemoveObjectOptions) -> Option { self.remove_object_inner(bucket_name, object_name, opts).await.err() } pub async fn remove_object_inner( &self, bucket_name: &str, object_name: &str, opts: RemoveObjectOptions, ) -> Result { let mut url_values = HashMap::new(); if opts.version_id != "" { url_values.insert("versionId".to_string(), opts.version_id.clone()); } let mut headers = HeaderMap::new(); if opts.governance_bypass { headers.insert(X_AMZ_BYPASS_GOVERNANCE_RETENTION, HeaderValue::from_static("true")); //amzBypassGovernance } let resp = self .execute_method( http::Method::DELETE, &mut RequestMetadata { bucket_name: bucket_name.to_string(), object_name: object_name.to_string(), content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), query_values: url_values, custom_header: headers, content_body: ReaderImpl::Body(Bytes::new()), content_length: 0, content_md5_base64: "".to_string(), stream_sha256: false, trailer: HeaderMap::new(), pre_sign_url: Default::default(), add_crc: Default::default(), extra_pre_sign_header: Default::default(), bucket_location: Default::default(), expires: Default::default(), }, ) .await?; Ok(RemoveObjectResult { object_name: object_name.to_string(), object_version_id: opts.version_id, delete_marker: resp.headers().get(X_AMZ_DELETE_MARKER).map_or(false, |v| v == "true"), delete_marker_version_id: resp .headers() .get(X_AMZ_VERSION_ID) .and_then(|v| v.to_str().ok()) .unwrap_or_default() .to_string(), ..Default::default() }) } pub async fn remove_objects_with_result( self: Arc, bucket_name: &str, objects_rx: Receiver, opts: RemoveObjectsOptions, ) -> Receiver { let (result_tx, result_rx) = mpsc::channel(1); let self_clone = Arc::clone(&self); let bucket_name_owned = bucket_name.to_string(); tokio::spawn(async move { self_clone .remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts) .await; }); result_rx } pub async fn remove_objects( self: Arc, bucket_name: &str, objects_rx: Receiver, opts: RemoveObjectsOptions, ) -> Receiver { let (error_tx, error_rx) = mpsc::channel(1); let self_clone = Arc::clone(&self); let bucket_name_owned = bucket_name.to_string(); let (result_tx, mut result_rx) = mpsc::channel(1); tokio::spawn(async move { self_clone .remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts) .await; }); tokio::spawn(async move { while let Some(res) = result_rx.recv().await { if res.err.is_none() { continue; } error_tx .send(RemoveObjectError { object_name: res.object_name, version_id: res.object_version_id, err: res.err, ..Default::default() }) .await; } }); error_rx } pub async fn remove_objects_inner( &self, bucket_name: &str, mut objects_rx: Receiver, result_tx: &Sender, opts: RemoveObjectsOptions, ) -> Result<(), std::io::Error> { let max_entries = 1000; let mut finish = false; let mut url_values = HashMap::new(); url_values.insert("delete".to_string(), "".to_string()); loop { if finish { break; } let mut count = 0; let mut batch = Vec::::new(); while let Some(object) = objects_rx.recv().await { if has_invalid_xml_char(&object.name) { let remove_result = self .remove_object_inner( bucket_name, &object.name, RemoveObjectOptions { version_id: object.version_id.map(|id| id.to_string()).unwrap_or_default(), governance_bypass: opts.governance_bypass, ..Default::default() }, ) .await?; let remove_result_clone = remove_result.clone(); if let Some(err) = &remove_result.err { match to_error_response(err).code { S3ErrorCode::InvalidArgument | S3ErrorCode::NoSuchVersion => { continue; } _ => (), } result_tx.send(remove_result_clone.clone()).await; } result_tx.send(remove_result_clone).await; continue; } batch.push(object); count += 1; if count >= max_entries { break; } } if count == 0 { break; } if count < max_entries { finish = true; } let mut headers = HeaderMap::new(); if opts.governance_bypass { headers.insert(X_AMZ_BYPASS_GOVERNANCE_RETENTION, HeaderValue::from_static("true")); } let remove_bytes = generate_remove_multi_objects_request(&batch); let resp = self .execute_method( http::Method::POST, &mut RequestMetadata { bucket_name: bucket_name.to_string(), query_values: url_values.clone(), content_body: ReaderImpl::Body(Bytes::from(remove_bytes.clone())), content_length: remove_bytes.len() as i64, content_md5_base64: base64_encode(&HashAlgorithm::Md5.hash_encode(&remove_bytes).as_ref()), content_sha256_hex: base64_encode(&HashAlgorithm::SHA256.hash_encode(&remove_bytes).as_ref()), custom_header: headers, object_name: "".to_string(), stream_sha256: false, trailer: HeaderMap::new(), pre_sign_url: Default::default(), add_crc: Default::default(), extra_pre_sign_header: Default::default(), bucket_location: Default::default(), expires: Default::default(), }, ) .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); } } process_remove_multi_objects_response(ReaderImpl::Body(Bytes::from(body_vec)), &batch, result_tx.clone()).await; } Ok(()) } pub async fn remove_incomplete_upload(&self, bucket_name: &str, object_name: &str) -> Result<(), std::io::Error> { let upload_ids = self.find_upload_ids(bucket_name, object_name)?; for upload_id in upload_ids { self.abort_multipart_upload(bucket_name, object_name, &upload_id).await?; } Ok(()) } pub async fn abort_multipart_upload( &self, bucket_name: &str, object_name: &str, upload_id: &str, ) -> Result<(), std::io::Error> { let mut url_values = HashMap::new(); url_values.insert("uploadId".to_string(), upload_id.to_string()); let resp = self .execute_method( http::Method::DELETE, &mut RequestMetadata { bucket_name: bucket_name.to_string(), object_name: object_name.to_string(), query_values: url_values, content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), custom_header: HeaderMap::new(), content_body: ReaderImpl::Body(Bytes::new()), content_length: 0, content_md5_base64: "".to_string(), stream_sha256: false, trailer: HeaderMap::new(), pre_sign_url: Default::default(), add_crc: Default::default(), extra_pre_sign_header: Default::default(), bucket_location: Default::default(), expires: Default::default(), }, ) .await?; let resp_status = resp.status(); let h = resp.headers().clone(); //if resp.is_some() { if resp.status() != StatusCode::NO_CONTENT { let error_response: ErrorResponse; match resp.status() { StatusCode::NOT_FOUND => { error_response = ErrorResponse { code: S3ErrorCode::NoSuchUpload, message: "The specified multipart upload does not exist.".to_string(), bucket_name: bucket_name.to_string(), key: object_name.to_string(), request_id: resp .headers() .get("x-amz-request-id") .and_then(|v| v.to_str().ok()) .unwrap_or_default() .to_string(), host_id: resp .headers() .get("x-amz-id-2") .and_then(|v| v.to_str().ok()) .unwrap_or_default() .to_string(), region: resp .headers() .get("x-amz-bucket-region") .and_then(|v| v.to_str().ok()) .unwrap_or_default() .to_string(), ..Default::default() }; } _ => { return Err(std::io::Error::other(http_resp_to_error_response( resp_status, &h, vec![], bucket_name, object_name, ))); } } return Err(std::io::Error::other(error_response)); } //} Ok(()) } } #[derive(Debug, Default)] #[allow(dead_code)] pub struct RemoveObjectError { object_name: String, #[allow(dead_code)] version_id: String, err: Option, } impl Display for RemoveObjectError { fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result { if let Some(err) = &self.err { write!(f, "{}", err.to_string()) } else { write!(f, "unexpected remove object error result") } } } #[derive(Debug, Default)] pub struct RemoveObjectResult { pub object_name: String, pub object_version_id: String, pub delete_marker: bool, pub delete_marker_version_id: String, pub err: Option, } impl Clone for RemoveObjectResult { fn clone(&self) -> Self { Self { object_name: self.object_name.clone(), object_version_id: self.object_version_id.clone(), delete_marker: self.delete_marker, delete_marker_version_id: self.delete_marker_version_id.clone(), err: None, //err } } } pub struct RemoveObjectsOptions { pub governance_bypass: bool, } pub fn generate_remove_multi_objects_request(objects: &[ObjectInfo]) -> Vec { let escape_xml = |value: &str| -> String { value .replace('&', "&") .replace('<', "<") .replace('>', ">") .replace('\"', """) .replace('\'', "'") }; let request: DeleteMultiObjects = DeleteMultiObjects { quiet: false, objects: objects .iter() .map(|object| DeleteObject { key: object.name.clone(), version_id: object.version_id.map(|v| v.to_string()).unwrap_or_default(), }) .collect(), }; match request.marshal_msg() { Ok(body) => body.into_bytes(), Err(_) => { let mut body = String::new(); body.push_str("false"); for object in objects { body.push_str(""); body.push_str(""); body.push_str(&escape_xml(&object.name)); body.push_str(""); if object.version_id.is_some() { body.push_str(""); body.push_str(&escape_xml(&object.version_id.as_ref().map(|v| v.to_string()).unwrap_or_default())); body.push_str(""); } body.push_str(""); } body.push_str(""); body.into_bytes() } } } pub async fn process_remove_multi_objects_response( body: ReaderImpl, objects: &[ObjectInfo], result_tx: Sender, ) { let mut body_vec = Vec::new(); match body { ReaderImpl::Body(content_body) => { body_vec = content_body.to_vec(); } ReaderImpl::ObjectBody(mut object_body) => match object_body.read_all().await { Ok(content) => { body_vec = content; } Err(err) => { for object in objects { let version_id = object.version_id.as_ref().map(|v| v.to_string()).unwrap_or_default(); let _ = result_tx .send(RemoveObjectResult { object_name: object.name.clone(), object_version_id: version_id, err: Some(std::io::Error::other(ErrorResponse { code: S3ErrorCode::Custom("ReadDeleteResponseFailed".into()), message: format!("read multi remove response failed: {err}"), bucket_name: object.bucket.clone(), key: object.name.clone(), resource: "".to_string(), request_id: "".to_string(), host_id: "".to_string(), region: "".to_string(), server: "".to_string(), status_code: StatusCode::OK, })), ..Default::default() }) .await; } return; } }, } #[derive(Debug, Deserialize)] #[serde(rename = "DeleteResult")] struct Deleted { #[serde(rename = "Deleted", default)] deleted: Vec, #[serde(rename = "Error", default)] error: Vec, } #[derive(Debug, Deserialize)] struct DeleteResultDeleted { #[serde(rename = "Key")] key: String, #[serde(rename = "VersionId", default)] version_id: String, #[serde(rename = "DeleteMarker")] deletemarker: bool, #[serde(rename = "DeleteMarkerVersionId", default)] deletemarker_version_id: String, } #[derive(Debug, Deserialize)] struct DeleteResultError { #[serde(rename = "Key")] key: String, #[serde(rename = "VersionId", default)] version_id: String, #[serde(rename = "Code")] code: String, #[serde(rename = "Message")] message: String, } let mut pending = HashSet::with_capacity(objects.len()); for object in objects { pending.insert((object.name.clone(), object.version_id.as_ref().map(|v| v.to_string()).unwrap_or_default())); } let body = String::from_utf8_lossy(&body_vec).into_owned(); let parsed: Deleted = match quick_xml::de::from_str(&body) { Ok(parsed) => parsed, Err(err) => { for object in objects { let version_id = object.version_id.as_ref().map(|v| v.to_string()).unwrap_or_default(); let _ = result_tx .send(RemoveObjectResult { object_name: object.name.clone(), object_version_id: version_id, err: Some(std::io::Error::other(ErrorResponse { code: S3ErrorCode::Custom("UnmarshalDeleteResponseFailed".into()), message: format!( "unmarshal multi remove response failed: {err}; response_body={}", body.chars().take(DELETE_RESPONSE_PREVIEW_LEN).collect::() ), bucket_name: object.bucket.clone(), key: object.name.clone(), resource: "".to_string(), request_id: "".to_string(), host_id: "".to_string(), region: "".to_string(), server: "".to_string(), status_code: StatusCode::OK, })), ..Default::default() }) .await; } return; } }; for deleted in parsed.deleted { if !pending.remove(&(deleted.key.clone(), deleted.version_id.clone())) { continue; } let _ = result_tx .send(RemoveObjectResult { object_name: deleted.key, object_version_id: deleted.version_id, delete_marker: deleted.deletemarker, delete_marker_version_id: deleted.deletemarker_version_id, err: None, }) .await; } for removed in parsed.error { if !pending.remove(&(removed.key.clone(), removed.version_id.clone())) { continue; } let _ = result_tx .send(RemoveObjectResult { object_name: removed.key.clone(), object_version_id: removed.version_id, err: Some(std::io::Error::other(ErrorResponse { code: S3ErrorCode::Custom(removed.code.into()), message: removed.message, bucket_name: "".to_string(), key: removed.key, resource: "".to_string(), request_id: "".to_string(), host_id: "".to_string(), region: "".to_string(), server: "".to_string(), status_code: StatusCode::OK, })), ..Default::default() }) .await; } for (object_name, object_version_id) in pending { let bucket_name = objects .iter() .find(|object| { object.name == object_name && object.version_id.as_ref().map(|v| v.to_string()) == Some(object_version_id.clone()) }) .map(|o| o.bucket.clone()) .unwrap_or_default(); let object_name = object_name; let object_version_id = object_version_id; let error_message = format!( "remove response did not contain an entry for object {} with version {}", object_name, object_version_id ); let _ = result_tx .send(RemoveObjectResult { object_name: object_name.clone(), object_version_id: object_version_id.clone(), err: Some(std::io::Error::other(ErrorResponse { code: S3ErrorCode::Custom("UnmatchedDeleteResponseEntry".into()), message: error_message, bucket_name, key: object_name, resource: "".to_string(), request_id: "".to_string(), host_id: "".to_string(), region: "".to_string(), server: "".to_string(), status_code: StatusCode::OK, })), ..Default::default() }) .await; } } fn has_invalid_xml_char(str: &str) -> bool { false }