From 47f365d385d00f18783767a32a155f7d1882f2e8 Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 30 Dec 2024 13:59:14 +0800 Subject: [PATCH] feat:#165 copyobject v1 --- ecstore/src/config/common.rs | 6 +- ecstore/src/erasure.rs | 7 +- ecstore/src/heal/data_usage_cache.rs | 4 +- ecstore/src/set_disk.rs | 273 +++++++++++++++++++-------- ecstore/src/sets.rs | 61 +++++- ecstore/src/store.rs | 111 ++++++++++- ecstore/src/store_api.rs | 165 +++++++++++----- iam/src/store/object.rs | 8 +- rustfs/src/storage/ecfs.rs | 127 ++++++++++++- rustfs/src/storage/options.rs | 69 ++++++- scripts/run.sh | 1 + 11 files changed, 671 insertions(+), 161 deletions(-) diff --git a/ecstore/src/config/common.rs b/ecstore/src/config/common.rs index f25f557cc..e52f827e5 100644 --- a/ecstore/src/config/common.rs +++ b/ecstore/src/config/common.rs @@ -6,7 +6,7 @@ use super::{storageclass, Config, GLOBAL_StorageClass, KVS}; use crate::config::error::is_not_found; use crate::disk::RUSTFS_META_BUCKET; use crate::error::{Error, Result}; -use crate::store_api::{HTTPRangeSpec, ObjectInfo, ObjectOptions, PutObjReader, StorageAPI}; +use crate::store_api::{ObjectInfo, ObjectOptions, PutObjReader, StorageAPI}; use crate::store_err::is_err_object_not_found; use crate::utils::path::SLASH_SEPARATOR; use http::HeaderMap; @@ -31,7 +31,6 @@ lazy_static! { } pub async fn read_config(api: Arc, file: &str) -> Result> { let (data, _obj) = read_config_with_metadata(api, file, &ObjectOptions::default()).await?; - Ok(data) } @@ -40,10 +39,9 @@ async fn read_config_with_metadata( file: &str, opts: &ObjectOptions, ) -> Result<(Vec, ObjectInfo)> { - let range = HTTPRangeSpec::nil(); let h = HeaderMap::new(); let mut rd = api - .get_object_reader(RUSTFS_META_BUCKET, file, range, h, opts) + .get_object_reader(RUSTFS_META_BUCKET, file, None, h, opts) .await .map_err(|err| { if is_err_object_not_found(&err) { diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index 64c04f523..4121ddf62 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -10,8 +10,8 @@ use std::fmt::Debug; use std::io::ErrorKind; use tokio::io::DuplexStream; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tracing::info; use tracing::warn; +use tracing::{error, info}; // use tracing::debug; use uuid::Uuid; @@ -263,7 +263,10 @@ impl Erasure { .await { Ok(n) => n, - Err(err) => return (bytes_writed, Some(err)), + Err(err) => { + error!("write_data_blocks err {:?}", &err); + return (bytes_writed, Some(err)); + } }; bytes_writed += writed_n; diff --git a/ecstore/src/heal/data_usage_cache.rs b/ecstore/src/heal/data_usage_cache.rs index fbc6003b6..ead588021 100644 --- a/ecstore/src/heal/data_usage_cache.rs +++ b/ecstore/src/heal/data_usage_cache.rs @@ -383,7 +383,7 @@ impl DataUsageCache { .get_object_reader( RUSTFS_META_BUCKET, path.to_str().unwrap(), - HTTPRangeSpec::nil(), + None, HeaderMap::new(), &ObjectOptions { no_lock: true, @@ -404,7 +404,7 @@ impl DataUsageCache { .get_object_reader( RUSTFS_META_BUCKET, name, - HTTPRangeSpec::nil(), + None, HeaderMap::new(), &ObjectOptions { no_lock: true, diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 79c78eb92..42d15a918 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -8,8 +8,6 @@ use std::{ use crate::config::error::is_not_found; use crate::global::GLOBAL_MRFState; -use crate::heal::heal_ops::{HealEntryFn, HealSequence}; -use crate::heal::mrf::PartialOperation; use crate::{ bitrot::{bitrot_verify, close_bitrot_writers, new_bitrot_filereader, new_bitrot_filewriter, BitrotFileWriter}, cache_value::metacache_set::{list_path_raw, ListPathRawOptions}, @@ -54,11 +52,16 @@ use crate::{ }, xhttp, }; +use crate::{disk::STORAGE_FORMAT_FILE, heal::mrf::PartialOperation}; use crate::{file_meta::file_info_from_raw, heal::data_usage_cache::DataUsageCache}; use crate::{ heal::data_scanner::{globalHealConfig, HEAL_DELETE_DANGLING}, store_api::ListObjectVersionsInfo, }; +use crate::{ + heal::heal_ops::{HealEntryFn, HealSequence}, + utils::path::path_join_buf, +}; use bytesize::ByteSize; use chrono::Utc; use futures::future::join_all; @@ -76,7 +79,7 @@ use rand::{ {seq::SliceRandom, Rng}, }; use reader::reader::EtagReader; -use s3s::dto::StreamingBlob; +use s3s::{dto::StreamingBlob, Body}; use sha2::{Digest, Sha256}; use std::hash::Hash; use std::time::SystemTime; @@ -574,10 +577,10 @@ impl SetDisks { bucket: &str, prefix: &str, files: &[FileInfo], - // write_quorum: usize, - ) -> Vec> { + write_quorum: usize, + ) -> Result<()> { let mut futures = Vec::with_capacity(disks.len()); - let mut errors = Vec::with_capacity(disks.len()); + let mut errs = Vec::with_capacity(disks.len()); for (i, disk) in disks.iter().enumerate() { let mut file_info = files[i].clone(); @@ -595,14 +598,42 @@ impl SetDisks { for result in results { match result { Ok(_) => { - errors.push(None); + errs.push(None); } Err(e) => { - errors.push(Some(e)); + errs.push(Some(e)); } } } - errors + + if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) { + // TODO: 并发 + for (i, err) in errs.iter().enumerate() { + if err.is_some() { + continue; + } + + if let Some(disk) = disks[i].as_ref() { + let _ = disk + .delete( + bucket, + &path_join_buf(&[prefix, STORAGE_FORMAT_FILE]), + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + .map_err(|e| { + warn!("write meta revert err {:?}", e); + e + }); + } + } + + return Err(err); + } + Ok(()) } fn get_upload_id_dir(bucket: &str, object: &str, upload_id: &str) -> String { @@ -1747,12 +1778,17 @@ impl SetDisks { } #[allow(clippy::too_many_arguments)] + #[tracing::instrument( + level = "debug", + skip( writer,disks,fi,files), + fields(start_time=?time::OffsetDateTime::now_utc()) + )] async fn get_object_with_fileinfo( // &self, bucket: &str, object: &str, - offset: i64, - length: i64, + offset: usize, + length: usize, writer: &mut DuplexStream, fi: FileInfo, files: Vec, @@ -1762,10 +1798,10 @@ impl SetDisks { ) -> Result<()> { let (disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, &files, &fi); - let total_size = fi.size as i64; + let total_size = fi.size; let length = { - if length < 0 { + if length == 0 { total_size - offset } else { length @@ -1791,13 +1827,13 @@ impl SetDisks { let (last_part_index, _) = fi.to_part_offset(end_offset)?; // debug!( - // "get_object_with_fileinfo end offset:{}, part_index:{},part_offset:{}", + // "get_object_with_fileinfo end offset:{}, last_part_index:{},part_offset:{}", // end_offset, last_part_index, 0 // ); let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); - let mut total_readed: i64 = 0; + let mut total_readed = 0; for i in part_index..=last_part_index { if total_readed == length { break; @@ -1805,12 +1841,12 @@ impl SetDisks { let part_number = fi.parts[i].number; let part_size = fi.parts[i].size; - let mut part_length = part_size - part_offset as usize; - if part_length > (length - total_readed) as usize { - part_length = (length - total_readed) as usize + let mut part_length = part_size - part_offset; + if part_length > length - total_readed { + part_length = length - total_readed } - let till_offset = erasure.shard_file_offset(part_offset.try_into().unwrap(), part_length, part_size); + let till_offset = erasure.shard_file_offset(part_offset, part_length, part_size); let mut readers = Vec::with_capacity(disks.len()); for (idx, disk_op) in disks.iter().enumerate() { // debug!("read part_path {}", &part_path); @@ -1844,13 +1880,12 @@ impl SetDisks { // "read part {} part_offset {},part_length {},part_size {} ", // part_number, part_offset, part_length, part_size // ); - let (written, mut err) = erasure - .decode(writer, readers, part_offset as usize, part_length, part_size) - .await; + let (written, mut err) = erasure.decode(writer, readers, part_offset, part_length, part_size).await; if let Some(e) = err.as_ref() { if written == part_length { match e.downcast_ref::() { Some(DiskError::FileNotFound) | Some(DiskError::FileCorrupt) => { + error!("erasure.decode err 111 {:?}", &e); GLOBAL_MRFState .add_partial(PartialOperation { bucket: bucket.to_string(), @@ -1870,11 +1905,12 @@ impl SetDisks { } } if let Some(err) = err { + error!("erasure.decode err {} {:?}", written, &err); return Err(err); } // debug!("ec decode {} writed size {}", part_number, n); - total_readed += part_length as i64; + total_readed += part_length; part_offset = 0; } @@ -3495,19 +3531,26 @@ impl SetDisks { #[async_trait::async_trait] impl ObjectIO for SetDisks { + #[tracing::instrument(level = "debug", skip(self))] async fn get_object_reader( &self, bucket: &str, object: &str, - _rs: HTTPRangeSpec, - _h: HeaderMap, + range: Option, + h: HeaderMap, opts: &ObjectOptions, ) -> Result { - let (fi, files, disks) = self.get_object_fileinfo(bucket, object, opts, true).await?; + let (fi, files, disks) = self + .get_object_fileinfo(bucket, object, opts, true) + .await + .map_err(|err| to_object_err(err, vec![bucket, object]))?; let object_info = fi.to_object_info(bucket, object, opts.versioned || opts.version_suspended); if object_info.delete_marker { - return Err(Error::new(DiskError::FileNotFound)); + if opts.version_id.is_none() { + return Err(to_object_err(Error::new(DiskError::FileNotFound), vec![bucket, object])); + } + return Err(to_object_err(Error::new(StorageError::MethodNotAllowed), vec![bucket, object])); } // if object_info.size == 0 { @@ -3519,17 +3562,28 @@ impl ObjectIO for SetDisks { // }); // } - let rs = HTTPRangeSpec::from_object_info(&object_info, opts.part_number); - let (offset, length) = rs.get_offset_length(object_info.size.try_into().unwrap())?; + if object_info.size == 0 { + if let Some(rs) = range { + let _ = rs.get_offset_length(object_info.size)?; + } - // debug!("get_object_reader offset:{}, length:{}", offset, length); + let reader = GetObjectReader { + stream: StreamingBlob::from(Body::from(Vec::new())), + object_info, + }; + return Ok(reader); + } + + // TODO: remote let (rd, mut wd) = tokio::io::duplex(fi.erasure.block_size); - // let disks = self.disks.read().await; + + let (reader, offset, length) = + GetObjectReader::new(StreamingBlob::wrap(tokio_util::io::ReaderStream::new(rd)), range, &object_info, opts, &h)?; // let disks = disks.clone(); - let bucket = String::from(bucket); - let object = String::from(object); + let bucket = bucket.to_owned(); + let object = object.to_owned(); let set_index = self.set_index; let pool_index = self.pool_index; tokio::spawn(async move { @@ -3542,14 +3596,6 @@ impl ObjectIO for SetDisks { }; }); - let read_stream = tokio_util::io::ReaderStream::new(rd); - - // let rd: Box = Box::new(rd); - - let reader = GetObjectReader { - stream: StreamingBlob::wrap(read_stream), - object_info, - }; Ok(reader) } @@ -3585,7 +3631,7 @@ impl ObjectIO for SetDisks { _ns = Some(ns_lock); } - let mut user_defined = opts.user_defined.clone(); + let mut user_defined = opts.user_defined.clone().unwrap_or_default(); let sc_parity_drives = { if let Some(sc) = GLOBAL_StorageClass.get() { @@ -3789,6 +3835,110 @@ impl StorageAPI for SetDisks { unimplemented!() } + async fn copy_object( + &self, + src_bucket: &str, + src_object: &str, + _dst_bucket: &str, + _dst_object: &str, + src_info: &mut ObjectInfo, + src_opts: &ObjectOptions, + dst_opts: &ObjectOptions, + ) -> Result { + // FIXME: TODO: + + if !src_info.metadata_only { + return Err(Error::new(StorageError::NotImplemented)); + } + + let disks = self.get_disks_internal().await; + + let (mut metas, errs) = { + if let Some(vid) = &src_opts.version_id { + Self::read_all_fileinfo(&disks, "", src_bucket, src_object, vid, true, false).await + } else { + Self::read_all_xl(&disks, src_bucket, src_object, true, false).await + } + }; + + let (read_quorum, write_quorum) = match Self::object_quorum_from_meta(&metas, &errs, self.default_parity_count) { + Ok((r, w)) => (r as usize, w as usize), + Err(mut err) => { + if ErasureError::ErasureReadQuorum.is(&err) + && !src_bucket.starts_with(RUSTFS_META_BUCKET) + && self + .delete_if_dang_ling(src_bucket, src_object, &metas, &errs, &HashMap::new(), src_opts.clone()) + .await + .is_ok() + { + if src_opts.version_id.is_some() { + err = Error::new(DiskError::FileVersionNotFound) + } else { + err = Error::new(DiskError::FileNotFound) + } + } + return Err(to_object_err(err, vec![src_bucket, src_object])); + } + }; + + let (online_disks, mod_time, etag) = Self::list_online_disks(&disks, &metas, &errs, read_quorum); + + let mut fi = Self::pick_valid_fileinfo(&metas, mod_time, etag, read_quorum) + .map_err(|e| to_object_err(e, vec![src_bucket, src_object]))?; + + if fi.deleted { + if src_opts.version_id.is_none() { + return Err(to_object_err(Error::new(DiskError::FileNotFound), vec![src_bucket, src_object])); + } + return Err(to_object_err(Error::new(StorageError::MethodNotAllowed), vec![src_bucket, src_object])); + } + + let version_id = { + if src_info.version_only { + if let Some(vid) = &dst_opts.version_id { + Some(Uuid::parse_str(vid)?) + } else { + Some(Uuid::new_v4()) + } + } else { + src_info.version_id + } + }; + + let inline_data = fi.inline_data(); + fi.metadata = src_info.user_defined.clone(); + + if let Some(ud) = src_info.user_defined.as_mut() { + if let Some(etag) = &src_info.etag { + ud.insert("etag".to_owned(), etag.clone()); + } + } + + let mod_time = OffsetDateTime::now_utc(); + + for fi in metas.iter_mut() { + if fi.is_valid() { + fi.metadata = src_info.user_defined.clone(); + fi.mod_time = Some(mod_time); + fi.version_id = version_id; + fi.versioned = src_opts.versioned || src_opts.version_suspended; + + if !fi.inline_data() { + fi.data = None; + } + + if inline_data { + fi.set_inline_data(); + } + } + } + + Self::write_unique_file_info(&online_disks, "", src_bucket, src_object, &metas, write_quorum) + .await + .map_err(|e| to_object_err(e, vec![src_bucket, src_object]))?; + + Ok(fi.to_object_info(src_bucket, src_object, src_opts.versioned || src_opts.version_suspended)) + } async fn delete_objects( &self, bucket: &str, @@ -4333,7 +4483,7 @@ impl StorageAPI for SetDisks { let disks = disks.clone(); - let mut user_defined = opts.user_defined.clone(); + let mut user_defined = opts.user_defined.clone().unwrap_or_default(); if let Some(ref etag) = opts.preserve_etag { user_defined.insert("etag".to_owned(), etag.clone()); @@ -4371,8 +4521,6 @@ impl StorageAPI for SetDisks { write_quorum += 1 } - let _ = write_quorum; - let mut fi = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives); fi.version_id = if let Some(vid) = &opts.version_id { @@ -4418,42 +4566,17 @@ impl StorageAPI for SetDisks { let upload_path = Self::get_upload_id_dir(bucket, object, upload_uuid.as_str()); - let errs = Self::write_unique_file_info( + Self::write_unique_file_info( &shuffle_disks, bucket, RUSTFS_META_MULTIPART_BUCKET, upload_path.as_str(), &parts_metadatas, + write_quorum, ) - .await; + .await + .map_err(|e| to_object_err(e, vec![bucket, object]))?; - if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) { - // TODO: 并发 - for (i, err) in errs.iter().enumerate() { - if err.is_some() { - continue; - } - - if let Some(disk) = shuffle_disks[i].as_ref() { - let _ = disk - .delete( - RUSTFS_META_MULTIPART_BUCKET, - upload_path.as_str(), - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .map_err(|e| { - warn!("write meta revert err {:?}", e); - e - }); - } - } - - return Err(err); - } // evalDisks Ok(MultipartUploadResult { upload_id }) @@ -4612,7 +4735,7 @@ impl StorageAPI for SetDisks { // etag let etag = { - if let Some(etag) = opts.user_defined.get("etag") { + if let Some(Some(etag)) = opts.user_defined.as_ref().map(|v| v.get("etag")) { etag.clone() } else { get_complete_multipart_md5(&uploaded_parts) diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index f5cd4949a..cbf83eea8 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -27,8 +27,9 @@ use crate::{ ListMultipartsInfo, ListObjectVersionsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartInfo, MultipartUploadResult, ObjectIO, ObjectInfo, ObjectOptions, ObjectToDelete, PartInfo, PutObjReader, StorageAPI, }, + store_err::StorageError, store_init::{check_format_erasure_values, get_format_erasure_in_quorum, load_format_erasure_all, save_format_file}, - utils::hash, + utils::{hash, path::path_join_buf}, }; use crate::heal::heal_ops::HealSequence; @@ -292,7 +293,7 @@ impl ObjectIO for Sets { &self, bucket: &str, object: &str, - range: HTTPRangeSpec, + range: Option, h: HeaderMap, opts: &ObjectOptions, ) -> Result { @@ -360,6 +361,62 @@ impl StorageAPI for Sets { async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { unimplemented!() } + async fn copy_object( + &self, + src_bucket: &str, + src_object: &str, + dst_bucket: &str, + dst_object: &str, + src_info: &mut ObjectInfo, + src_opts: &ObjectOptions, + dst_opts: &ObjectOptions, + ) -> Result { + let src_set = self.get_disks_by_key(src_object); + let dst_set = self.get_disks_by_key(dst_object); + + let cp_src_dst_same = path_join_buf(&[src_bucket, src_object]) == path_join_buf(&[dst_bucket, dst_object]); + + if cp_src_dst_same { + if let (Some(src_vid), Some(dst_vid)) = (&src_opts.version_id, &dst_opts.version_id) { + if src_vid == dst_vid { + return src_set + .copy_object(src_bucket, src_object, dst_bucket, dst_object, src_info, src_opts, dst_opts) + .await; + } + } + + if !dst_opts.versioned && src_opts.version_id.is_none() { + return src_set + .copy_object(src_bucket, src_object, dst_bucket, dst_object, src_info, src_opts, dst_opts) + .await; + } + + if dst_opts.versioned && src_opts.version_id != dst_opts.version_id { + src_info.version_only = true; + return src_set + .copy_object(src_bucket, src_object, dst_bucket, dst_object, src_info, src_opts, dst_opts) + .await; + } + } + + let put_opts = ObjectOptions { + user_defined: dst_opts.user_defined.clone(), + versioned: dst_opts.versioned, + version_id: dst_opts.version_id.clone(), + mod_time: dst_opts.mod_time, + ..Default::default() + }; + + if let Some(put_object_reader) = src_info.put_object_reader.as_mut() { + return dst_set.put_object(dst_bucket, dst_object, put_object_reader, &put_opts).await; + } + + Err(Error::new(StorageError::InvalidArgument( + src_bucket.to_owned(), + src_object.to_owned(), + "put_object_reader2 is none".to_owned(), + ))) + } async fn delete_objects( &self, bucket: &str, diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 9014eaa35..843c81d31 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -24,7 +24,7 @@ use crate::store_err::{ }; use crate::store_init::ec_drives_no_config; use crate::utils::crypto::base64_decode; -use crate::utils::path::{decode_dir_object, encode_dir_object, SLASH_SEPARATOR}; +use crate::utils::path::{decode_dir_object, encode_dir_object, path_join_buf, SLASH_SEPARATOR}; use crate::utils::xml; use crate::{ bucket::metadata::BucketMetadata, @@ -116,7 +116,7 @@ impl ECStore { ) .await; - debug!("endpoint_pools: {:?}", endpoint_pools); + // debug!("endpoint_pools: {:?}", endpoint_pools); let mut common_parity_drives = 0; @@ -533,6 +533,39 @@ impl ECStore { Ok(idx) } + async fn get_pool_idx_no_lock(&self, bucket: &str, object: &str, size: i64) -> Result { + let idx = match self.get_pool_idx_existing_no_lock(bucket, object).await { + Ok(res) => res, + Err(err) => { + if !is_err_object_not_found(&err) { + return Err(err); + } + + if let Some(idx) = self.get_available_pool_idx(bucket, object, size).await { + idx + } else { + return Err(to_object_err(Error::new(DiskError::DiskFull), vec![bucket, object])); + } + } + }; + + Ok(idx) + } + + async fn get_pool_idx_existing_no_lock(&self, bucket: &str, object: &str) -> Result { + self.get_pool_idx_existing_with_opts( + bucket, + object, + &ObjectOptions { + no_lock: true, + skip_decommissioned: true, + skip_rebalancing: true, + ..Default::default() + }, + ) + .await + } + async fn get_pool_idx_existing_with_opts(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { let (pinfo, _) = self.get_pool_info_existing_with_opts(bucket, object, opts).await?; Ok(pinfo.index) @@ -1052,7 +1085,7 @@ impl ObjectIO for ECStore { &self, bucket: &str, object: &str, - range: HTTPRangeSpec, + range: Option, h: HeaderMap, opts: &ObjectOptions, ) -> Result { @@ -1289,6 +1322,76 @@ impl StorageAPI for ECStore { Ok(info) } + // TODO: review + async fn copy_object( + &self, + src_bucket: &str, + src_object: &str, + dst_bucket: &str, + dst_object: &str, + src_info: &mut ObjectInfo, + src_opts: &ObjectOptions, + dst_opts: &ObjectOptions, + ) -> Result { + check_copy_obj_args(src_bucket, src_object)?; + check_copy_obj_args(dst_bucket, dst_object)?; + + let src_object = utils::path::encode_dir_object(src_object); + let dst_object = utils::path::encode_dir_object(dst_object); + + let cp_src_dst_same = path_join_buf(&[src_bucket, &src_object]) == path_join_buf(&[dst_bucket, &dst_object]); + + // TODO: nslock + + let pool_idx = self + .get_pool_idx_no_lock(src_bucket, &src_object, src_info.size as i64) + .await?; + + if cp_src_dst_same { + if let (Some(src_vid), Some(dst_vid)) = (&src_opts.version_id, &dst_opts.version_id) { + if src_vid == dst_vid { + return self.pools[pool_idx] + .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, dst_opts) + .await; + } + } + + if !dst_opts.versioned && src_opts.version_id.is_none() { + return self.pools[pool_idx] + .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, dst_opts) + .await; + } + + if dst_opts.versioned && src_opts.version_id != dst_opts.version_id { + src_info.version_only = true; + return self.pools[pool_idx] + .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, dst_opts) + .await; + } + } + + let put_opts = ObjectOptions { + user_defined: src_info.user_defined.clone(), + versioned: dst_opts.versioned, + version_id: dst_opts.version_id.clone(), + no_lock: true, + mod_time: dst_opts.mod_time, + ..Default::default() + }; + + if let Some(put_object_reader) = src_info.put_object_reader.as_mut() { + return self.pools[pool_idx] + .put_object(dst_bucket, &dst_object, put_object_reader, &put_opts) + .await; + } + + Err(Error::new(StorageError::InvalidArgument( + src_bucket.to_owned(), + src_object.to_owned(), + "put_object_reader is none".to_owned(), + ))) + } + // TODO: review async fn delete_objects( &self, @@ -2183,7 +2286,7 @@ fn check_object_name_for_length_and_slash(bucket: &str, object: &str) -> Result< Ok(()) } -fn _check_copy_obj_args(bucket: &str, object: &str) -> Result<()> { +fn check_copy_obj_args(bucket: &str, object: &str) -> Result<()> { check_bucket_and_object_names(bucket, object) } diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 725e6de7b..986a90e31 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -8,7 +8,7 @@ use crate::{ xhttp, }; use futures::StreamExt; -use http::HeaderMap; +use http::{HeaderMap, HeaderValue}; use madmin::heal_commands::HealResultItem; use rmp_serde::Serializer; use s3s::{dto::StreamingBlob, Body}; @@ -233,7 +233,7 @@ impl FileInfo { } // to_part_offset 取offset 所在的part index, 返回part index, offset - pub fn to_part_offset(&self, offset: i64) -> Result<(usize, i64)> { + pub fn to_part_offset(&self, offset: usize) -> Result<(usize, usize)> { if offset == 0 { return Ok((0, 0)); } @@ -241,11 +241,11 @@ impl FileInfo { let mut part_offset = offset; for (i, part) in self.parts.iter().enumerate() { let part_index = i; - if part_offset < part.size as i64 { + if part_offset < part.size { return Ok((part_index, part_offset)); } - part_offset -= part.size as i64 + part_offset -= part.size } Err(Error::msg("part not found")) @@ -442,9 +442,44 @@ pub struct GetObjectReader { } impl GetObjectReader { - // pub fn new(stream: StreamingBlob, object_info: ObjectInfo) -> Self { - // GetObjectReader { stream, object_info } - // } + #[tracing::instrument(level = "debug", skip(reader))] + pub fn new( + reader: StreamingBlob, + rs: Option, + oi: &ObjectInfo, + opts: &ObjectOptions, + _h: &HeaderMap, + ) -> Result<(Self, usize, usize)> { + let mut rs = rs; + + if let Some(part_number) = opts.part_number { + if rs.is_none() { + rs = HTTPRangeSpec::from_object_info(oi, part_number); + } + } + + if let Some(rs) = rs { + let (off, length) = rs.get_offset_length(oi.size)?; + + return Ok(( + GetObjectReader { + stream: reader, + object_info: oi.clone(), + }, + off, + length, + )); + } else { + return Ok(( + GetObjectReader { + stream: reader, + object_info: oi.clone(), + }, + 0, + oi.size, + )); + } + } pub async fn read_all(&mut self) -> Result> { let mut data = Vec::new(); @@ -463,61 +498,47 @@ impl GetObjectReader { #[derive(Debug)] pub struct HTTPRangeSpec { pub is_suffix_length: bool, - pub start: i64, - pub end: i64, + pub start: usize, + pub end: Option, } impl HTTPRangeSpec { - pub fn nil() -> Self { - Self { - is_suffix_length: false, - start: -1, - end: -1, - } - } - - pub fn is_nil(&self) -> bool { - self.start == -1 && self.end == -1 - } - pub fn from_object_info(oi: &ObjectInfo, part_number: usize) -> Self { - let mut l = oi.parts.len(); - if part_number < l { - l = part_number; + pub fn from_object_info(oi: &ObjectInfo, part_number: usize) -> Option { + if oi.size == 0 || oi.parts.is_empty() { + return None; } let mut start = 0; let mut end = -1; - for i in 0..l { + for i in 0..oi.parts.len().min(part_number) { start = end + 1; end = start + oi.parts[i].size as i64 - 1 } - HTTPRangeSpec { + Some(HTTPRangeSpec { is_suffix_length: false, - start, - end, - } + start: start as usize, + end: { + if end < 0 { + None + } else { + Some(end as usize) + } + }, + }) } - pub fn get_offset_length(&self, res_size: i64) -> Result<(i64, i64)> { - if self.start == 0 && self.end == 0 { - return Ok((0, res_size)); - } - + pub fn get_offset_length(&self, res_size: usize) -> Result<(usize, usize)> { let len = self.get_length(res_size)?; let mut start = self.start; if self.is_suffix_length { - start = self.start + res_size + start = res_size - self.start } Ok((start, len)) } - pub fn get_length(&self, res_size: i64) -> Result { - if self.is_nil() { - return Ok(res_size); - } - + pub fn get_length(&self, res_size: usize) -> Result { if self.is_suffix_length { - let specified_len = -self.start; // 假设 h.start 是一个 i64 类型 + let specified_len = self.start; // 假设 h.start 是一个 i64 类型 let mut range_length = specified_len; if specified_len > res_size { @@ -527,21 +548,21 @@ impl HTTPRangeSpec { return Ok(range_length); } - if self.start > res_size { + if self.start >= res_size { return Err(Error::msg("The requested range is not satisfiable")); } - if self.end > -1 { - let mut end = self.end; + if let Some(end) = self.end { + let mut end = end; if res_size <= end { end = res_size - 1; } - let range_length = end - self.start - 1; + let range_length = end - self.start + 1; return Ok(range_length); } - if self.end == -1 { + if self.end.is_none() { let range_length = res_size - self.start; return Ok(range_length); } @@ -555,7 +576,7 @@ pub struct ObjectOptions { // Use the maximum parity (N/2), used when saving server configuration files pub max_parity: bool, pub mod_time: Option, - pub part_number: usize, + pub part_number: Option, pub delete_prefix: bool, pub version_id: Option, @@ -569,7 +590,7 @@ pub struct ObjectOptions { pub data_movement: bool, pub src_pool_idx: usize, - pub user_defined: HashMap, + pub user_defined: Option>, pub preserve_etag: Option, pub metadata_chg: bool, @@ -631,7 +652,7 @@ impl From for CompletePart { } } -#[derive(Debug, Default, Clone)] +#[derive(Debug, Default)] pub struct ObjectInfo { pub bucket: String, pub name: String, @@ -652,9 +673,41 @@ pub struct ObjectInfo { pub content_encoding: Option, pub num_versions: usize, pub successor_mod_time: Option, - // pub put_object_reader: Option, + pub put_object_reader: Option, pub etag: Option, pub inlined: bool, + pub metadata_only: bool, + pub version_only: bool, +} + +impl Clone for ObjectInfo { + fn clone(&self) -> Self { + Self { + bucket: self.bucket.clone(), + name: self.name.clone(), + mod_time: self.mod_time, + size: self.size, + actual_size: self.actual_size, + is_dir: self.is_dir, + user_defined: self.user_defined.clone(), + parity_blocks: self.parity_blocks, + data_blocks: self.data_blocks, + version_id: self.version_id, + delete_marker: self.delete_marker, + user_tags: self.user_tags.clone(), + parts: self.parts.clone(), + is_latest: self.is_latest, + content_type: self.content_type.clone(), + content_encoding: self.content_encoding.clone(), + num_versions: self.num_versions, + successor_mod_time: self.successor_mod_time, + put_object_reader: None, // reader can not clone + etag: self.etag.clone(), + inlined: self.inlined, + metadata_only: self.metadata_only, + version_only: self.version_only, + } + } } impl ObjectInfo { @@ -839,7 +892,7 @@ pub trait ObjectIO: Send + Sync + 'static { &self, bucket: &str, object: &str, - range: HTTPRangeSpec, + range: Option, h: HeaderMap, opts: &ObjectOptions, ) -> Result; @@ -889,6 +942,16 @@ pub trait StorageAPI: ObjectIO { async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result; // PutObject ObjectIO // CopyObject + async fn copy_object( + &self, + src_bucket: &str, + src_object: &str, + dst_bucket: &str, + dst_object: &str, + src_info: &mut ObjectInfo, + src_opts: &ObjectOptions, + dst_opts: &ObjectOptions, + ) -> Result; async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result; async fn delete_objects( &self, diff --git a/iam/src/store/object.rs b/iam/src/store/object.rs index 15a4abfbe..8199903c5 100644 --- a/iam/src/store/object.rs +++ b/iam/src/store/object.rs @@ -117,13 +117,7 @@ impl Store for ObjectStore { debug!("load iam config, path: {}", path.as_ref()); let mut reader = self .object_api - .get_object_reader( - Self::BUCKET_NAME, - path.as_ref(), - HTTPRangeSpec::nil(), - Default::default(), - &Default::default(), - ) + .get_object_reader(Self::BUCKET_NAME, path.as_ref(), None, Default::default(), &Default::default()) .await .map_err(crate::Error::EcstoreError)?; diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index a3b08e1f0..422a1d2c8 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -30,6 +30,7 @@ use ecstore::store_api::ObjectOptions; use ecstore::store_api::ObjectToDelete; use ecstore::store_api::PutObjReader; use ecstore::store_api::StorageAPI; +use ecstore::utils::path::path_join_buf; use ecstore::utils::xml; use ecstore::xhttp; use futures::pin_mut; @@ -52,7 +53,9 @@ use transform_stream::AsyncTryStream; use uuid::Uuid; use crate::storage::error::to_s3_error; -use crate::storage::options::extract_metadata_from_mime; +use crate::storage::options::copy_dst_opts; +use crate::storage::options::copy_src_opts; +use crate::storage::options::{extract_metadata_from_mime, get_opts}; macro_rules! try_ { ($result:expr) => { @@ -119,13 +122,87 @@ impl S3 for FS { #[tracing::instrument(level = "debug", skip(self, req))] async fn copy_object(&self, req: S3Request) -> S3Result> { - let input = req.input; - let (_bucket, _key) = match input.copy_source { + let CopyObjectInput { + copy_source, + bucket, + key, + .. + } = req.input; + let (src_bucket, src_key, version_id) = match copy_source { CopySource::AccessPoint { .. } => return Err(s3_error!(NotImplemented)), - CopySource::Bucket { ref bucket, ref key, .. } => (bucket, key), + CopySource::Bucket { + ref bucket, + ref key, + version_id, + } => (bucket.to_string(), key.to_string(), version_id.map(|v| v.to_string())), }; - let output = CopyObjectOutput { ..Default::default() }; + // warn!("copy_object {}/{}, to {}/{}", &src_bucket, &src_key, &bucket, &key); + + let mut src_opts = copy_src_opts(&src_bucket, &src_key, &req.headers).map_err(to_s3_error)?; + + src_opts.version_id = version_id.clone(); + + let mut get_opts = ObjectOptions { + version_id: src_opts.version_id.clone(), + versioned: src_opts.versioned, + version_suspended: src_opts.version_suspended, + ..Default::default() + }; + + let dst_opts = copy_dst_opts(&bucket, &key, version_id, &req.headers, None) + .await + .map_err(to_s3_error)?; + + let cp_src_dst_same = path_join_buf(&[&src_bucket, &src_key]) == path_join_buf(&[&bucket, &key]); + + if cp_src_dst_same { + get_opts.no_lock = true; + } + + let Some(store) = new_object_layer_fn() else { + return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); + }; + + let h = HeaderMap::new(); + + let gr = store + .get_object_reader(&src_bucket, &src_key, None, h, &get_opts) + .await + .map_err(to_s3_error)?; + + let mut src_info = gr.object_info.clone(); + + if cp_src_dst_same { + src_info.metadata_only = true; + } + + src_info.put_object_reader = Some(PutObjReader { + stream: gr.stream, + content_length: gr.object_info.size as usize, + }); + + // check quota + // TODO: src metadada + // TODO: src tags + + let oi = store + .copy_object(&src_bucket, &src_key, &bucket, &key, &mut src_info, &src_opts, &dst_opts) + .await + .map_err(to_s3_error)?; + + // warn!("copy_object oi {:?}", &oi); + + let copy_object_result = CopyObjectResult { + e_tag: oi.etag, + last_modified: oi.mod_time.map(Timestamp::from), + ..Default::default() + }; + + let output = CopyObjectOutput { + copy_object_result: Some(copy_object_result), + ..Default::default() + }; Ok(S3Response::new(output)) } @@ -312,16 +389,46 @@ impl S3 for FS { // warn!("get_object input {:?}, vid {:?}", &req.input, req.input.version_id); let GetObjectInput { - bucket, key, version_id, .. + bucket, + key, + version_id, + part_number, + range, + .. } = req.input; - let range = HTTPRangeSpec::nil(); + // let range = HTTPRangeSpec::nil(); let h = HeaderMap::new(); - let metadata = extract_metadata(&req.headers); + let part_number = part_number.map(|v| v as usize); - let opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, Some(metadata)) + if let Some(part_num) = part_number { + if part_num == 0 { + return Err(s3_error!(InvalidArgument, "part_numer invalid")); + } + } + + let rs = range.map(|v| match v { + Range::Int { first, last } => HTTPRangeSpec { + is_suffix_length: false, + start: first as usize, + end: last.map(|v| v as usize), + }, + Range::Suffix { length } => HTTPRangeSpec { + is_suffix_length: true, + start: length as usize, + end: None, + }, + }); + + if rs.is_some() && part_number.is_some() { + return Err(s3_error!(InvalidArgument, "range and part_number invalid")); + } + + // let metadata = extract_metadata(&req.headers); + + let opts: ObjectOptions = get_opts(&bucket, &key, version_id, part_number, &req.headers) .await .map_err(to_s3_error)?; @@ -330,7 +437,7 @@ impl S3 for FS { }; let reader = store - .get_object_reader(bucket.as_str(), key.as_str(), range, h, &opts) + .get_object_reader(bucket.as_str(), key.as_str(), rs, h, &opts) .await .map_err(to_s3_error)?; diff --git a/rustfs/src/storage/options.rs b/rustfs/src/storage/options.rs index 9ae3cd89e..847d67791 100644 --- a/rustfs/src/storage/options.rs +++ b/rustfs/src/storage/options.rs @@ -56,6 +56,55 @@ pub async fn del_opts( Ok(opts) } +pub async fn get_opts( + bucket: &str, + object: &str, + vid: Option, + part_num: Option, + headers: &HeaderMap, +) -> Result { + let versioned = BucketVersioningSys::prefix_enabled(bucket, object).await; + let version_suspended = BucketVersioningSys::prefix_suspended(bucket, object).await; + + let vid = vid.map(|v| v.as_str().trim().to_owned()); + + if let Some(ref id) = vid { + if let Err(_err) = Uuid::parse_str(id.as_str()) { + return Err(Error::new(StorageError::InvalidVersionID( + bucket.to_owned(), + object.to_owned(), + id.clone(), + ))); + } + + if !versioned { + return Err(Error::new(StorageError::InvalidArgument( + bucket.to_owned(), + object.to_owned(), + id.clone(), + ))); + } + } + + let mut opts = get_default_opts(headers, None, false) + .map_err(|err| Error::new(StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), err.to_string())))?; + + opts.version_id = { + if is_dir_object(object) && vid.is_none() { + Some(Uuid::nil().to_string()) + } else { + vid + } + }; + + opts.part_number = part_num; + + opts.version_suspended = version_suspended; + opts.versioned = versioned; + + Ok(opts) +} + pub async fn put_opts( bucket: &str, object: &str, @@ -102,18 +151,30 @@ pub async fn put_opts( Ok(opts) } +pub async fn copy_dst_opts( + bucket: &str, + object: &str, + vid: Option, + headers: &HeaderMap, + metadata: Option>, +) -> Result { + put_opts(bucket, object, vid, headers, metadata).await +} + +pub fn copy_src_opts(_bucket: &str, _object: &str, headers: &HeaderMap) -> Result { + get_default_opts(headers, None, false) +} + pub fn put_opts_from_headers( headers: &HeaderMap, metadata: Option>, ) -> Result { - let metadata = metadata.unwrap_or_default(); - get_default_opts(headers, metadata, false) } -fn get_default_opts( +pub fn get_default_opts( _headers: &HeaderMap, - metadata: HashMap, + metadata: Option>, _copy_source: bool, ) -> Result { Ok(ObjectOptions { diff --git a/scripts/run.sh b/scripts/run.sh index 5b1326e6a..55b3f7ffa 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -18,6 +18,7 @@ fi export RUSTFS_STORAGE_CLASS_INLINE_BLOCK="512 KB" + DATA_DIR_ARG="./target/volume/test{0...4}" # DATA_DIR_ARG="./target/volume/test"