diff --git a/Cargo.lock b/Cargo.lock index df15b64e6..a7a81c8ad 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1563,6 +1563,7 @@ dependencies = [ "tracing-error", "tracing-subscriber", "transform-stream", + "uuid", ] [[package]] @@ -2242,9 +2243,9 @@ checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" [[package]] name = "uuid" -version = "1.9.1" +version = "1.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5de17fd2f7da591098415cff336e12965a28061ddace43b59cb3c430179c9439" +checksum = "81dfa00651efa65069b0b6b651f4aaa31ba9e3c3ce0137aaad053604ee7e0314" dependencies = [ "getrandom", "rand", diff --git a/Cargo.toml b/Cargo.toml index ca8430b09..88dc7fa8b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,7 +4,7 @@ members = ["rustfs", "ecstore", "e2e_test", "common/protos"] [workspace.package] edition = "2021" -license = "MIT OR Apache-2.0" +license = "Apache-2.0" repository = "https://github.com/rustfs/rustfs" rust-version = "1.75" version = "0.0.1" diff --git a/TODO.md b/TODO.md index 973bef0ac..f19a90c3d 100644 --- a/TODO.md +++ b/TODO.md @@ -29,16 +29,27 @@ - [x] 提交完成 CompleteMultipartUpload - [x] 取消上传 AbortMultipartUpload - [x] 下载 GetObject - - [ ] 删除 DeleteObjects + - [x] 删除 DeleteObjects - [ ] 版本控制 - [ ] 对象锁 - [ ] 复制 CopyObject - [ ] 详情 HeadObject - + - [ ] 对象预先签名(get、put、head、post) + ## 扩展功能 -- [ ] 版本控制 -- [ ] 对象锁 -- [ ] 修复 +- [ ] 用户管理 +- [ ] Policy管理 +- [ ] AK/SK分配管理 +- [ ] data scanner统计和对象修复 - [ ] 桶配额 - [ ] 桶只读 +- [ ] 桶复制 +- [ ] 桶事件通知 +- [ ] 桶公开、桶私有 +- [ ] 对象生命周期管理 +- [ ] prometheus对接 +- [ ] 日志收集和日志外发 +- [ ] 对象压缩 +- [ ] STS +- [ ] 分层(阿里云、腾讯云、S3远程对接) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index aae55759f..a58d5fe7f 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -1,7 +1,7 @@ use super::{endpoint::Endpoint, error::DiskError, format::FormatV3}; use super::{ - DeleteOptions, DiskAPI, FileReader, FileWriter, MetaCacheEntry, ReadMultipleReq, ReadMultipleResp, ReadOptions, - RenameDataResp, VolumeInfo, WalkDirOptions, + DeleteOptions, DiskAPI, FileInfoVersions, FileReader, FileWriter, MetaCacheEntry, ReadMultipleReq, ReadMultipleResp, + ReadOptions, RenameDataResp, VolumeInfo, WalkDirOptions, }; use crate::disk::{LocalFileReader, LocalFileWriter, STORAGE_FORMAT_FILE}; use crate::{ @@ -138,8 +138,34 @@ impl LocalDisk { Ok(()) } + pub async fn move_to_trash(&self, delete_path: &PathBuf, _recursive: bool, _immediate_purge: bool) -> Result<()> { + let trash_path = self.get_object_path(super::RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?; + // TODO: 清空回收站 + if let Err(err) = fs::rename(&delete_path, &trash_path).await { + match err.kind() { + ErrorKind::NotFound => (), + _ => { + warn!("delete_file rename {:?} err {:?}", &delete_path, &err); + return Err(Error::from(err)); + } + } + } + + // FIXME: 先清空回收站吧,有时间再添加判断逻辑 + let _ = fs::remove_dir_all(&trash_path).await; + + // TODO: immediate + Ok(()) + } + // #[tracing::instrument(skip(self))] - pub async fn delete_file(&self, base_path: &PathBuf, delete_path: &PathBuf, recursive: bool, _immediate: bool) -> Result<()> { + pub async fn delete_file( + &self, + base_path: &PathBuf, + delete_path: &PathBuf, + recursive: bool, + immediate_purge: bool, + ) -> Result<()> { debug!("delete_file {:?}\n base_path:{:?}", &delete_path, &base_path); if is_root_path(base_path) || is_root_path(delete_path) { @@ -153,29 +179,7 @@ impl LocalDisk { } if recursive { - let trash_path = self.get_object_path(super::RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?; - - if let Some(dir_path) = trash_path.parent() { - fs::create_dir_all(dir_path).await?; - } - - debug!("delete_file ranme to trash {:?} to {:?}", &delete_path, &trash_path); - - // TODO: 清空回收站 - if let Err(err) = fs::rename(&delete_path, &trash_path).await { - match err.kind() { - ErrorKind::NotFound => (), - _ => { - warn!("delete_file rename {:?} err {:?}", &delete_path, &err); - return Err(Error::from(err)); - } - } - } - - // FIXME: 先清空回收站吧,有时间再添加判断逻辑 - let _ = fs::remove_dir_all(&trash_path).await; - - // TODO: immediate + self.move_to_trash(delete_path, recursive, immediate_purge).await?; } else { if delete_path.is_dir() { if let Err(err) = fs::remove_dir(&delete_path).await { @@ -253,6 +257,52 @@ impl LocalDisk { Ok((data, modtime)) } + + async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &Vec) -> Result<()> { + let volume_dir = self.get_bucket_path(volume)?; + let xlpath = self.get_object_path(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str())?; + + let (data, _) = match self.read_all_data(volume, volume_dir.as_path(), &xlpath).await { + Ok(res) => res, + Err(_err) => { + // TODO: check if not found return err + + (Vec::new(), OffsetDateTime::UNIX_EPOCH) + } + }; + + if data.is_empty() { + return Err(Error::new(DiskError::FileNotFound)); + } + + let mut fm = FileMeta::default(); + + fm.unmarshal_msg(&data)?; + + for fi in fis { + let data_dir = fm.delete_version(fi)?; + + if data_dir.is_some() { + let dir_path = self.get_object_path(volume, format!("{}/{}", path, data_dir.unwrap().to_string()).as_str())?; + + self.move_to_trash(&dir_path, true, false).await?; + } + } + + // 没有版本了,删除xl.meta + if fm.versions.is_empty() { + self.delete_file(&volume_dir, &xlpath, true, false).await?; + return Ok(()); + } + + // 更新xl.meta + let buf = fm.marshal_msg()?; + + self.write_all(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str(), buf) + .await?; + + Ok(()) + } } fn is_root_path(path: impl AsRef) -> bool { @@ -606,7 +656,7 @@ impl DiskAPI for LocalDisk { let (src_data_path, dst_data_path) = { let mut data_dir = String::new(); if !fi.is_remote() { - data_dir = utils::path::retain_slash(fi.data_dir.to_string().as_str()); + data_dir = utils::path::retain_slash(fi.data_dir.unwrap_or(Uuid::nil()).to_string().as_str()); } if !data_dir.is_empty() { @@ -645,10 +695,9 @@ impl DiskAPI for LocalDisk { let old_data_dir = meta .find_version(fi.version_id) .map(|(_, version)| { - version.get_data_dir().filter(|data_dir| { - warn!("get data dir {}", &data_dir); - meta.shard_data_dir_count(&fi.version_id, data_dir) == 0 - }) + version + .get_data_dir() + .filter(|data_dir| meta.shard_data_dir_count(&fi.version_id, &Some(data_dir.clone())) == 0) }) .unwrap_or_default(); @@ -819,6 +868,27 @@ impl DiskAPI for LocalDisk { Ok(RawFileInfo { buf }) } + async fn delete_versions( + &self, + volume: &str, + versions: Vec, + _opts: DeleteOptions, + ) -> Result>> { + let mut errs = Vec::with_capacity(versions.len()); + for _ in 0..versions.len() { + errs.push(None); + } + + for (i, ver) in versions.iter().enumerate() { + if let Err(e) = self.delete_versions_internal(volume, ver.name.as_str(), &ver.versions).await { + errs[i] = Some(e); + } else { + errs[i] = None; + } + } + + Ok(errs) + } async fn read_multiple(&self, req: ReadMultipleReq) -> Result> { let mut results = Vec::new(); let mut found = 0; diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 2fc1cf3c3..1256dbc3d 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -14,7 +14,7 @@ const STORAGE_FORMAT_FILE: &str = "xl.meta"; use crate::{ erasure::{ReadAt, Write}, - error::Result, + error::{Error, Result}, file_meta::FileMeta, store_api::{FileInfo, RawFileInfo}, }; @@ -85,9 +85,31 @@ pub trait DiskAPI: Debug + Send + Sync + 'static { opts: &ReadOptions, ) -> Result; async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result; + async fn delete_versions( + &self, + volume: &str, + versions: Vec, + opts: DeleteOptions, + ) -> Result>>; async fn read_multiple(&self, req: ReadMultipleReq) -> Result>; } +#[derive(Debug, Default, Clone)] +pub struct FileInfoVersions { + // Name of the volume. + pub volume: String, + + // Name of the file. + pub name: String, + + // Represents the latest mod time of the + // latest version. + pub latest_mod_time: Option, + + pub versions: Vec, + pub free_versions: Vec, +} + #[derive(Debug, Default, Clone, Serialize, Deserialize)] pub struct WalkDirOptions { // Bucket to scanner diff --git a/ecstore/src/disk/remote.rs b/ecstore/src/disk/remote.rs index dc1e99964..68398ac81 100644 --- a/ecstore/src/disk/remote.rs +++ b/ecstore/src/disk/remote.rs @@ -20,13 +20,14 @@ use tracing::info; use uuid::Uuid; use crate::{ - error::Result, + error::{Error, Result}, store_api::{FileInfo, RawFileInfo}, }; use super::{ - endpoint::Endpoint, DeleteOptions, DiskAPI, DiskOption, FileReader, FileWriter, MetaCacheEntry, ReadMultipleReq, - ReadMultipleResp, ReadOptions, RemoteFileReader, RemoteFileWriter, RenameDataResp, VolumeInfo, WalkDirOptions, + endpoint::Endpoint, DeleteOptions, DiskAPI, DiskOption, FileInfoVersions, FileReader, FileWriter, MetaCacheEntry, + ReadMultipleReq, ReadMultipleResp, ReadOptions, RemoteFileReader, RemoteFileWriter, RenameDataResp, VolumeInfo, + WalkDirOptions, }; #[derive(Debug)] @@ -40,36 +41,28 @@ impl RemoteDisk { pub async fn new(ep: &Endpoint, _opt: &DiskOption) -> Result { let root = fs::canonicalize(ep.url.path()).await?; - Ok(Self { + Ok(Self { channel: RwLock::new(None), url: ep.url.clone(), - root + root, }) } fn get_client(&self) -> NodeServiceClient> { let channel = { let read_lock = self.channel.read().unwrap(); - + if let Some(ref channel) = *read_lock { channel.clone() } else { - let addr = format!( - "{}://{}:{}", - self.url.scheme(), - self.url.host_str().unwrap(), - self.url.port().unwrap() - ); + let addr = format!("{}://{}:{}", self.url.scheme(), self.url.host_str().unwrap(), self.url.port().unwrap()); info!("disk url: {:?}", addr); let connector = tonic_Endpoint::from_shared(addr).unwrap(); - - let new_channel = tokio::runtime::Runtime::new() - .unwrap() - .block_on(connector.connect()) - .unwrap(); - + + let new_channel = tokio::runtime::Runtime::new().unwrap().block_on(connector.connect()).unwrap(); + *self.channel.write().unwrap() = Some(new_channel.clone()); - + new_channel } }; @@ -347,6 +340,15 @@ impl DiskAPI for RemoteDisk { Ok(raw_file_info) } + async fn delete_versions( + &self, + _volume: &str, + _versions: Vec, + _opts: DeleteOptions, + ) -> Result>> { + unimplemented!() + } + async fn read_multiple(&self, req: ReadMultipleReq) -> Result> { let read_multiple_req = serde_json::to_string(&req)?; let mut client = self.get_client(); diff --git a/ecstore/src/endpoints.rs b/ecstore/src/endpoints.rs index 6f781704d..9af5b963b 100644 --- a/ecstore/src/endpoints.rs +++ b/ecstore/src/endpoints.rs @@ -408,6 +408,11 @@ impl AsMut> for EndpointServerPools { } impl EndpointServerPools { + pub fn from_volumes(server_addr: &str, endpoints: Vec) -> Result<(EndpointServerPools, SetupType)> { + let layouts = DisksLayout::try_from(endpoints.as_slice())?; + + Self::create_server_endpoints(server_addr, &layouts) + } /// validates and creates new endpoints from input args, supports /// both ellipses and without ellipses transparently. pub fn create_server_endpoints(server_addr: &str, disks_layout: &DisksLayout) -> Result<(EndpointServerPools, SetupType)> { diff --git a/ecstore/src/lib.rs b/ecstore/src/lib.rs index beec2d9ec..85e1c9017 100644 --- a/ecstore/src/lib.rs +++ b/ecstore/src/lib.rs @@ -1,8 +1,8 @@ mod bucket_meta; mod chunk_stream; pub mod disk; -mod disks_layout; -mod endpoints; +pub mod disks_layout; +pub mod endpoints; pub mod erasure; pub mod error; mod file_meta; diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index b9adb1f9d..63815089d 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -1,3 +1,6 @@ +use std::collections::HashMap; + +use futures::future::join_all; use http::HeaderMap; use uuid::Uuid; @@ -7,11 +10,11 @@ use crate::{ DiskStore, }, endpoints::PoolEndpoints, - error::Result, + error::{Error, Result}, set_disk::SetDisks, store_api::{ - BucketInfo, BucketOptions, CompletePart, GetObjectReader, HTTPRangeSpec, ListObjectsV2Info, MakeBucketOptions, - MultipartUploadResult, ObjectInfo, ObjectOptions, PartInfo, PutObjReader, StorageAPI, + BucketInfo, BucketOptions, CompletePart, DeletedObject, GetObjectReader, HTTPRangeSpec, ListObjectsV2Info, + MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, ObjectToDelete, PartInfo, PutObjReader, StorageAPI, }, utils::hash, }; @@ -110,6 +113,22 @@ impl Sets { // ) -> Vec> { // unimplemented!() // } + + async fn delete_prefix(&self, bucket: &str, object: &str) -> Result<()> { + let mut futures = Vec::new(); + let opt = ObjectOptions { + delete_prefix: true, + ..Default::default() + }; + + for set in self.disk_set.iter() { + futures.push(set.delete_object(bucket, object, opt.clone())); + } + + let _results = join_all(futures).await; + + Ok(()) + } } // #[derive(Debug)] @@ -122,6 +141,12 @@ impl Sets { // pub default_parity_count: usize, // } +struct DelObj { + // set_idx: usize, + orig_idx: usize, + obj: ObjectToDelete, +} + #[async_trait::async_trait] impl StorageAPI for Sets { async fn list_bucket(&self, _opts: &BucketOptions) -> Result> { @@ -134,7 +159,77 @@ impl StorageAPI for Sets { async fn get_bucket_info(&self, _bucket: &str, _opts: &BucketOptions) -> Result { unimplemented!() } + async fn delete_objects( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + ) -> Result<(Vec, Vec>)> { + // 默认返回值 + let mut del_objects = vec![DeletedObject::default(); objects.len()]; + let mut del_errs = Vec::with_capacity(objects.len()); + for _ in 0..objects.len() { + del_errs.push(None) + } + + let mut set_obj_map = HashMap::new(); + + // hash key + let mut i = 0; + for obj in objects.iter() { + let idx = self.get_hashed_set_index(obj.object_name.as_str()); + + if !set_obj_map.contains_key(&idx) { + set_obj_map.insert( + idx, + vec![DelObj { + // set_idx: idx, + orig_idx: i, + obj: obj.clone(), + }], + ); + } else { + if let Some(val) = set_obj_map.get_mut(&idx) { + val.push(DelObj { + // set_idx: idx, + orig_idx: i, + obj: obj.clone(), + }); + } + } + + i += 1; + } + + // TODO: 并发 + for (k, v) in set_obj_map { + let disks = self.get_disks(k); + let objs: Vec = v.iter().map(|v| v.obj.clone()).collect(); + let (dobjects, errs) = disks.delete_objects(bucket, objs, opts.clone()).await?; + + let mut i = 0; + for err in errs { + let obj = v.get(i).unwrap(); + + del_errs[obj.orig_idx] = err; + + del_objects[obj.orig_idx] = dobjects.get(i).unwrap().clone(); + + i += 1; + } + } + + Ok((del_objects, del_errs)) + } + async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { + if opts.delete_prefix { + self.delete_prefix(bucket, object).await?; + return Ok(ObjectInfo::default()); + } + + self.get_disks_by_key(object).delete_object(bucket, object, opts).await + } async fn list_objects_v2( &self, _bucket: &str, diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index b31a5fc4d..b0d42bbe6 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -1,24 +1,45 @@ use crate::{ bucket_meta::BucketMetadata, disk::{error::DiskError, DeleteOptions, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET}, - disks_layout::DisksLayout, endpoints::EndpointServerPools, error::{Error, Result}, peer::S3PeerSys, sets::Sets, store_api::{ - BucketInfo, BucketOptions, CompletePart, GetObjectReader, HTTPRangeSpec, ListObjectsInfo, ListObjectsV2Info, - MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, PartInfo, PutObjReader, StorageAPI, + BucketInfo, BucketOptions, CompletePart, DeletedObject, GetObjectReader, HTTPRangeSpec, ListObjectsInfo, + ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, ObjectToDelete, PartInfo, + PutObjReader, StorageAPI, }, store_init, utils, }; use futures::future::join_all; use http::HeaderMap; use s3s::{dto::StreamingBlob, Body}; -use std::collections::{HashMap, HashSet}; +use std::{ + collections::{HashMap, HashSet}, + sync::Arc, +}; +use time::OffsetDateTime; +use tokio::sync::Mutex; use tracing::{debug, warn}; use uuid::Uuid; +use lazy_static::lazy_static; + +lazy_static! { + pub static ref GLOBAL_OBJECT_API: Arc>> = Arc::new(Mutex::new(None)); +} + +pub fn new_object_layer_fn() -> Arc>> { + // 这里不需要显式地锁定和解锁,因为 Arc 提供了必要的线程安全性 + GLOBAL_OBJECT_API.clone() +} + +async fn set_object_layer(o: ECStore) { + let mut global_object_api = GLOBAL_OBJECT_API.lock().await; + *global_object_api = Some(o); +} + #[derive(Debug)] pub struct ECStore { pub id: uuid::Uuid, @@ -30,12 +51,12 @@ pub struct ECStore { } impl ECStore { - pub async fn new(address: String, endpoints: Vec) -> Result { - let layouts = DisksLayout::try_from(endpoints.as_slice())?; + pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<()> { + // let layouts = DisksLayout::try_from(endpoints.as_slice())?; let mut deployment_id = None; - let (endpoint_pools, _) = EndpointServerPools::create_server_endpoints(address.as_str(), &layouts)?; + // let (endpoint_pools, _) = EndpointServerPools::create_server_endpoints(address.as_str(), &layouts)?; let mut pools = Vec::with_capacity(endpoint_pools.as_ref().len()); let mut disk_map = HashMap::with_capacity(endpoint_pools.as_ref().len()); @@ -95,13 +116,17 @@ impl ECStore { let peer_sys = S3PeerSys::new(&endpoint_pools, local_disks.clone()); - Ok(ECStore { + let ec = ECStore { id: deployment_id.unwrap(), disk_map, pools, local_disks, peer_sys, - }) + }; + + set_object_layer(ec).await; + + Ok(()) } pub fn local_disks(&self) -> Vec { @@ -231,6 +256,77 @@ impl ECStore { Ok(()) } + async fn delete_prefix(&self, _bucket: &str, _object: &str) -> Result<()> { + unimplemented!() + } + + async fn get_pool_info_existing_with_opts( + &self, + bucket: &str, + object: &str, + opts: &ObjectOptions, + ) -> Result<(PoolObjInfo, Vec)> { + let mut futures = Vec::new(); + + for pool in self.pools.iter() { + futures.push(pool.get_object_info(bucket, object, opts)); + } + + let results = join_all(futures).await; + + let mut ress = Vec::new(); + + let mut i = 0; + + // join_all结果跟输入顺序一致 + for res in results { + let index = i; + + match res { + Ok(r) => { + ress.push(PoolObjInfo { + index, + object_info: r, + err: None, + }); + } + Err(e) => { + ress.push(PoolObjInfo { + index, + err: Some(e), + ..Default::default() + }); + } + } + i += 1; + } + + ress.sort_by(|a, b| { + let at = a.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH); + let bt = b.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH); + + at.cmp(&bt) + }); + + for res in ress { + // check + if res.err.is_none() { + // TODO: let errs = self.poolsWithObject() + return Ok((res, Vec::new())); + } + } + + let ret = PoolObjInfo::default(); + + Ok((ret, Vec::new())) + } +} + +#[derive(Debug, Default)] +pub struct PoolObjInfo { + pub index: usize, + pub object_info: ObjectInfo, + pub err: Option, } #[derive(Debug, Default)] @@ -305,7 +401,149 @@ impl StorageAPI for ECStore { Ok(info) } + async fn delete_objects( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + ) -> Result<(Vec, Vec>)> { + // encode object name + let objects: Vec = objects + .iter() + .map(|v| { + let mut v = v.clone(); + v.object_name = utils::path::encode_dir_object(v.object_name.as_str()); + v + }) + .collect(); + // 默认返回值 + let mut del_objects = vec![DeletedObject::default(); objects.len()]; + + let mut del_errs = Vec::with_capacity(objects.len()); + for _ in 0..objects.len() { + del_errs.push(None) + } + + // TODO: limte 限制并发数量 + let opt = ObjectOptions::default(); + // 取所有poolObjInfo + let mut futures = Vec::new(); + for obj in objects.iter() { + futures.push(self.get_pool_info_existing_with_opts(bucket, &obj.object_name, &opt)); + } + + let results = join_all(futures).await; + + // 记录pool Index 对应的objects pool_idx -> objects idx + let mut pool_index_objects = HashMap::new(); + + let mut i = 0; + for res in results { + match res { + Ok((pinfo, _)) => { + if pinfo.object_info.delete_marker && opts.version_id.is_empty() { + del_objects[i] = DeletedObject { + delete_marker: pinfo.object_info.delete_marker, + delete_marker_version_id: pinfo.object_info.version_id.map(|v| v.to_string()), + object_name: utils::path::decode_dir_object(&pinfo.object_info.name), + delete_marker_mtime: pinfo.object_info.mod_time, + ..Default::default() + }; + } + + if !pool_index_objects.contains_key(&pinfo.index) { + pool_index_objects.insert(pinfo.index, vec![i]); + } else { + // let mut vals = pool_index_objects. + if let Some(val) = pool_index_objects.get_mut(&pinfo.index) { + val.push(i); + } + } + } + Err(e) => { + //TODO: check not found + + del_errs[i] = Some(e) + } + } + + i += 1; + } + + if !pool_index_objects.is_empty() { + for sets in self.pools.iter() { + // 取pool idx 对应的 objects index + let vals = pool_index_objects.get(&sets.pool_idx); + if vals.is_none() { + continue; + } + + let obj_idxs = vals.unwrap(); + // 取对应obj,理论上不会none + let objs: Vec = obj_idxs + .iter() + .filter_map(|&idx| { + if let Some(obj) = objects.get(idx) { + Some(obj.clone()) + } else { + None + } + }) + .collect(); + + if objs.is_empty() { + continue; + } + + let (pdel_objs, perrs) = sets.delete_objects(bucket, objs, opts.clone()).await?; + + // perrs的顺序理论上跟obj_idxs顺序一致 + let mut i = 0; + for err in perrs { + let obj_idx = obj_idxs[i]; + + if err.is_some() { + del_errs[obj_idx] = err; + } + + let mut dobj = pdel_objs.get(i).unwrap().clone(); + dobj.object_name = utils::path::decode_dir_object(&dobj.object_name); + + del_objects[obj_idx] = dobj; + + i += 1; + } + } + } + + Ok((del_objects, del_errs)) + } + async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { + if opts.delete_prefix { + self.delete_prefix(bucket, &object).await?; + return Ok(ObjectInfo::default()); + } + + let object = utils::path::encode_dir_object(object); + let object = object.as_str(); + + // 查询在哪个pool + let (mut pinfo, errs) = self.get_pool_info_existing_with_opts(bucket, object, &opts).await?; + if pinfo.object_info.delete_marker && opts.version_id.is_empty() { + pinfo.object_info.name = utils::path::decode_dir_object(object); + return Ok(pinfo.object_info); + } + + if !errs.is_empty() { + // TODO: deleteObjectFromAllPools + } + + let mut obj = self.pools[pinfo.index].delete_object(bucket, object, opts.clone()).await?; + obj.name = utils::path::decode_dir_object(object); + + Ok(obj) + } async fn list_objects_v2( &self, bucket: &str, diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 09993833c..2cbe8005f 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -10,15 +10,15 @@ pub const ERASURE_ALGORITHM: &str = "rs-vandermonde"; pub const BLOCK_SIZE_V2: usize = 1048576; // 1M // #[derive(Debug, Clone)] -#[derive(Serialize, Deserialize, Debug, PartialEq, Clone)] +#[derive(Serialize, Deserialize, Debug, PartialEq, Clone, Default)] pub struct FileInfo { pub name: String, pub volume: String, - pub version_id: Uuid, + pub version_id: Option, pub erasure: ErasureInfo, pub deleted: bool, // DataDir of the file - pub data_dir: Uuid, + pub data_dir: Option, pub mod_time: Option, pub size: usize, pub data: Option>, @@ -27,7 +27,66 @@ pub struct FileInfo { pub is_latest: bool, } +// impl Default for FileInfo { +// fn default() -> Self { +// Self { +// version_id: Default::default(), +// erasure: Default::default(), +// deleted: Default::default(), +// data_dir: Default::default(), +// mod_time: None, +// size: Default::default(), +// data: Default::default(), +// fresh: Default::default(), +// name: Default::default(), +// volume: Default::default(), +// parts: Default::default(), +// is_latest: Default::default(), +// } +// } +// } + impl FileInfo { + pub fn new(object: &str, data_blocks: usize, parity_blocks: usize) -> Self { + let indexs = { + let cardinality = data_blocks + parity_blocks; + let mut nums = vec![0; cardinality]; + let key_crc = crc32fast::hash(object.as_bytes()); + + let start = key_crc as usize % cardinality; + for i in 1..=cardinality { + nums[i - 1] = 1 + ((start + i) % cardinality); + } + + nums + }; + Self { + erasure: ErasureInfo { + algorithm: String::from(ERASURE_ALGORITHM), + data_blocks, + parity_blocks, + block_size: BLOCK_SIZE_V2, + distribution: indexs, + ..Default::default() + }, + ..Default::default() + } + } + + pub fn is_valid(&self) -> bool { + if self.deleted { + return true; + } + + let data_blocks = self.erasure.data_blocks; + let parity_blocks = self.erasure.parity_blocks; + + (data_blocks >= parity_blocks) + && (data_blocks > 0) + && (self.erasure.index > 0 + && self.erasure.index <= data_blocks + parity_blocks + && self.erasure.distribution.len() == (data_blocks + parity_blocks)) + } pub fn is_remote(&self) -> bool { // TODO: when lifecycle false @@ -86,7 +145,7 @@ impl FileInfo { parity_blocks: self.erasure.parity_blocks, data_blocks: self.erasure.data_blocks, version_id: self.version_id, - deleted: self.deleted, + delete_marker: self.deleted, mod_time: self.mod_time, size: self.size, parts: self.parts.clone(), @@ -113,68 +172,6 @@ impl FileInfo { } } -impl Default for FileInfo { - fn default() -> Self { - Self { - version_id: Uuid::nil(), - erasure: Default::default(), - deleted: Default::default(), - data_dir: Uuid::nil(), - mod_time: None, - size: Default::default(), - data: Default::default(), - fresh: Default::default(), - name: Default::default(), - volume: Default::default(), - parts: Default::default(), - is_latest: Default::default(), - } - } -} - -impl FileInfo { - pub fn new(object: &str, data_blocks: usize, parity_blocks: usize) -> Self { - let indexs = { - let cardinality = data_blocks + parity_blocks; - let mut nums = vec![0; cardinality]; - let key_crc = crc32fast::hash(object.as_bytes()); - - let start = key_crc as usize % cardinality; - for i in 1..=cardinality { - nums[i - 1] = 1 + ((start + i) % cardinality); - } - - nums - }; - Self { - erasure: ErasureInfo { - algorithm: String::from(ERASURE_ALGORITHM), - data_blocks, - parity_blocks, - block_size: BLOCK_SIZE_V2, - distribution: indexs, - ..Default::default() - }, - ..Default::default() - } - } - - pub fn is_valid(&self) -> bool { - if self.deleted { - return true; - } - - let data_blocks = self.erasure.data_blocks; - let parity_blocks = self.erasure.parity_blocks; - - (data_blocks >= parity_blocks) - && (data_blocks > 0) - && (self.erasure.index > 0 - && self.erasure.index <= data_blocks + parity_blocks - && self.erasure.distribution.len() == (data_blocks + parity_blocks)) - } -} - #[derive(Serialize, Deserialize, Debug, PartialEq, Clone, Default)] pub struct ObjectPartInfo { // pub etag: Option, @@ -360,12 +357,15 @@ impl HTTPRangeSpec { } } -#[derive(Debug, Default)] +#[derive(Debug, Default, Clone)] 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 delete_prefix: bool, + pub version_id: String, } // impl Default for ObjectOptions { @@ -416,13 +416,13 @@ impl From for CompletePart { pub struct ObjectInfo { pub bucket: String, pub name: String, + pub mod_time: Option, + pub size: usize, pub is_dir: bool, pub parity_blocks: usize, pub data_blocks: usize, - pub version_id: Uuid, - pub deleted: bool, - pub mod_time: Option, - pub size: usize, + pub version_id: Option, + pub delete_marker: bool, pub parts: Vec, pub is_latest: bool, } @@ -471,13 +471,36 @@ pub struct ListObjectsV2Info { pub prefixes: Vec, } +#[derive(Debug, Default, Clone)] +pub struct ObjectToDelete { + pub object_name: String, + pub version_id: Option, +} +#[derive(Debug, Default, Clone)] +pub struct DeletedObject { + pub delete_marker: bool, + pub delete_marker_version_id: Option, + pub object_name: String, + pub version_id: Option, + // MTime of DeleteMarker on source that needs to be propagated to replica + pub delete_marker_mtime: Option, + // to support delete marker replication + // pub replication_state: ReplicationState, +} + #[async_trait::async_trait] pub trait StorageAPI { async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()>; async fn delete_bucket(&self, bucket: &str) -> Result<()>; async fn list_bucket(&self, opts: &BucketOptions) -> Result>; async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result; - + async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result; + async fn delete_objects( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + ) -> Result<(Vec, Vec>)>; async fn list_objects_v2( &self, bucket: &str, diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 89bd92a7c..d2d1a989c 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -43,7 +43,26 @@ tower.workspace = true tracing-error.workspace = true tracing-subscriber.workspace = true transform-stream.workspace = true +uuid = "1.10.0" [build-dependencies] prost-build.workspace = true -tonic-build.workspace = true \ No newline at end of file +tonic-build.workspace = true +http.workspace = true +bytes.workspace = true +futures.workspace = true +futures-util.workspace = true +# uuid = { version = "1.8.0", features = ["v4", "fast-rng", "serde"] } +ecstore = { path = "../ecstore" } +s3s = "0.10.0" +clap = { version = "4.5.7", features = ["derive"] } +tracing-subscriber = { version = "0.3.18", features = ["env-filter", "time"] } +hyper-util = { version = "0.1.5", features = [ + "tokio", + "server-auto", + "server-graceful", +] } +mime = "0.3.17" +transform-stream = "0.3.0" +netif = "0.1.6" +# pin-utils = "0.1.0" diff --git a/rustfs/src/grpc.rs b/rustfs/src/grpc.rs index 544f5810e..b9edd6a19 100644 --- a/rustfs/src/grpc.rs +++ b/rustfs/src/grpc.rs @@ -1,5 +1,6 @@ use ecstore::{ disk::{DeleteOptions, DiskStore, ReadMultipleReq, ReadOptions, WalkDirOptions}, + endpoints::EndpointServerPools, erasure::{ReadAt, Write}, peer::{LocalPeerS3Client, PeerS3Client}, store_api::{BucketOptions, FileInfo, MakeBucketOptions}, @@ -28,7 +29,9 @@ struct NodeService { pub local_peer: LocalPeerS3Client, } -pub fn make_server(local_disks: Vec) -> NodeServer { +pub fn make_server(endpoint_pools: EndpointServerPools) -> NodeServer { + // TODO: 参考rustfs创建 https://github.com/rustfs/s3-rustfs/issues/30#issuecomment-2339664516 + let local_disks = Vec::new(); let local_peer = LocalPeerS3Client::new(local_disks, None, None); NodeServer::new(NodeService { local_peer }) } diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index c0567d5e1..827df794f 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -4,7 +4,7 @@ mod service; mod storage; use clap::Parser; -use ecstore::{error::Result, store::ECStore}; +use ecstore::{endpoints::EndpointServerPools, error::Result, store::ECStore}; use grpc::make_server; use hyper_util::{ rt::{TokioExecutor, TokioIo}, @@ -15,7 +15,7 @@ use s3s::{auth::SimpleAuth, service::S3ServiceBuilder}; use service::hybrid; use std::{io::IsTerminal, net::SocketAddr, str::FromStr}; use tokio::net::TcpListener; -use tracing::{debug, info}; +use tracing::{debug, info, warn}; use tracing_error::ErrorLayer; use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt}; @@ -36,10 +36,13 @@ fn setup_tracing() { } fn main() -> Result<()> { + //解析获得到的参数 let opt = config::Opt::parse(); + //设置trace setup_tracing(); + //运行参数 run(opt) } @@ -47,7 +50,9 @@ fn main() -> Result<()> { async fn run(opt: config::Opt) -> Result<()> { debug!("opt: {:?}", &opt); + //监听地址,端口从参数中获取 let listener = TcpListener::bind(opt.address.clone()).await?; + //获取监听地址 let local_addr: SocketAddr = listener.local_addr()?; // let mut domain_name = { @@ -66,13 +71,15 @@ async fn run(opt: config::Opt) -> Result<()> { // }) // }; - let store: ECStore = ECStore::new(opt.address.clone(), opt.volumes.clone()).await?; - let local_disks = store.local_disks(); - + // 用于rpc + let (endpoint_pools, _) = EndpointServerPools::from_volumes(opt.address.clone().as_str(), opt.volumes.clone())?; // Setup S3 service + // 本项目使用s3s库来实现s3服务 let service = { - let mut b = S3ServiceBuilder::new(storage::ecfs::FS::new(store).await?); - + // let mut b = S3ServiceBuilder::new(storage::ecfs::FS::new(opt.address.clone(), endpoint_pools).await?); + let mut b = S3ServiceBuilder::new(storage::ecfs::FS::new()); + //设置AK和SK + //其中部份内容从config配置文件中读取 let mut access_key = String::from_str(config::DEFAULT_ACCESS_KEY).unwrap(); let mut secret_key = String::from_str(config::DEFAULT_SECRET_KEY).unwrap(); @@ -81,7 +88,7 @@ async fn run(opt: config::Opt) -> Result<()> { access_key = ak; secret_key = sk; } - + //显示info信息 info!("authentication is enabled {}, {}", &access_key, &secret_key); b.set_auth(SimpleAuth::from_single(access_key, secret_key)); @@ -101,47 +108,58 @@ async fn run(opt: config::Opt) -> Result<()> { b.build() }; + let rpc_service = make_server(endpoint_pools.clone()); - let hyper_service = service.into_shared(); + tokio::spawn(async move { + let hyper_service = service.into_shared(); - let hybrid_service = TowerToHyperService::new(hybrid(hyper_service, make_server(local_disks))); + let hybrid_service = TowerToHyperService::new(hybrid(hyper_service, rpc_service)); - let http_server = ConnBuilder::new(TokioExecutor::new()); - let graceful = hyper_util::server::graceful::GracefulShutdown::new(); + let http_server = ConnBuilder::new(TokioExecutor::new()); + let mut ctrl_c = std::pin::pin!(tokio::signal::ctrl_c()); + let graceful = hyper_util::server::graceful::GracefulShutdown::new(); + info!("server is running at http://{local_addr}"); - let mut ctrl_c = std::pin::pin!(tokio::signal::ctrl_c()); - - info!("server is running at http://{local_addr}"); - - loop { - let (socket, _) = tokio::select! { - res = listener.accept() => { - match res { - Ok(conn) => conn, - Err(err) => { - tracing::error!("error accepting connection: {err}"); - continue; + loop { + let (socket, _) = tokio::select! { + res = listener.accept() => { + match res { + Ok(conn) => conn, + Err(err) => { + tracing::error!("error accepting connection: {err}"); + continue; + } } } - } - _ = ctrl_c.as_mut() => { - break; - } - }; + _ = ctrl_c.as_mut() => { + break; + } + }; - let conn = http_server.serve_connection(TokioIo::new(socket), hybrid_service.clone()); - let conn = graceful.watch(conn.into_owned()); - tokio::spawn(async move { - let _ = conn.await; - }); - } + let conn = http_server.serve_connection(TokioIo::new(socket), hybrid_service.clone()); + let conn = graceful.watch(conn.into_owned()); + tokio::spawn(async move { + let _ = conn.await; + }); + } + + tokio::select! { + () = graceful.shutdown() => { + tracing::debug!("Gracefully shutdown!"); + }, + () = tokio::time::sleep(std::time::Duration::from_secs(10)) => { + tracing::debug!("Waited 10 seconds for graceful shutdown, aborting..."); + } + } + }); + + warn!(" init store"); + // init store + ECStore::new(opt.address.clone(), endpoint_pools.clone()).await?; tokio::select! { - () = graceful.shutdown() => { - tracing::debug!("Gracefully shutdown!"); - }, - () = tokio::time::sleep(std::time::Duration::from_secs(10)) => { - tracing::debug!("Waited 10 seconds for graceful shutdown, aborting..."); + _ = tokio::signal::ctrl_c() => { + } } diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 76b091aca..bfb9edd9c 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -1,11 +1,13 @@ use bytes::Bytes; use ecstore::disk::error::DiskError; +use ecstore::store::new_object_layer_fn; use ecstore::store_api::BucketOptions; use ecstore::store_api::CompletePart; use ecstore::store_api::HTTPRangeSpec; use ecstore::store_api::MakeBucketOptions; use ecstore::store_api::MultipartUploadResult; use ecstore::store_api::ObjectOptions; +use ecstore::store_api::ObjectToDelete; use ecstore::store_api::PutObjReader; use ecstore::store_api::StorageAPI; use futures::pin_mut; @@ -21,9 +23,9 @@ use s3s::{S3Request, S3Response}; use std::fmt::Debug; use std::str::FromStr; use transform_stream::AsyncTryStream; +use uuid::Uuid; use ecstore::error::Result; -use ecstore::store::ECStore; use tracing::debug; macro_rules! try_ { @@ -39,12 +41,13 @@ macro_rules! try_ { #[derive(Debug)] pub struct FS { - pub store: ECStore, + // pub store: ECStore, } impl FS { - pub async fn new(store: ECStore) -> Result { - Ok(Self { store }) + pub fn new() -> Self { + // let store: ECStore = ECStore::new(address, endpoint_pools).await?; + Self {} } } #[async_trait::async_trait] @@ -57,8 +60,15 @@ impl S3 for FS { async fn create_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + try_!( - self.store + store .make_bucket(&input.bucket, &MakeBucketOptions { force_create: true }) .await ); @@ -83,24 +93,137 @@ impl S3 for FS { async fn delete_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; - try_!(self.store.delete_bucket(&input.bucket).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + try_!(store.delete_bucket(&input.bucket).await); Ok(S3Response::new(DeleteBucketOutput {})) } #[tracing::instrument(level = "debug", skip(self, req))] async fn delete_object(&self, req: S3Request) -> S3Result> { - let _input = req.input; + let DeleteObjectInput { + bucket, key, version_id, .. + } = req.input; - let output = DeleteObjectOutput::default(); + let version_id = version_id + .as_ref() + .map(|v| match Uuid::parse_str(v) { + Ok(id) => Some(id), + Err(_) => None, + }) + .unwrap_or_default(); + let dobj = ObjectToDelete { + object_name: key, + version_id, + }; + + let objects: Vec = vec![dobj]; + + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + let (dobjs, _errs) = try_!(store.delete_objects(&bucket, objects, ObjectOptions::default()).await); + + // TODO: let errors; + + let (delete_marker, version_id) = { + if let Some((a, b)) = dobjs + .iter() + .map(|v| { + let delete_marker = { + if v.delete_marker { + Some(true) + } else { + None + } + }; + + let version_id = v.version_id.clone(); + + (delete_marker, version_id) + }) + .next() + { + (a, b) + } else { + (None, None) + } + }; + + let output = DeleteObjectOutput { + delete_marker, + version_id, + ..Default::default() + }; Ok(S3Response::new(output)) } #[tracing::instrument(level = "debug", skip(self, req))] async fn delete_objects(&self, req: S3Request) -> S3Result> { - let _input = req.input; + // info!("delete_objects args {:?}", req.input); - let output = DeleteObjectsOutput { ..Default::default() }; + let DeleteObjectsInput { bucket, delete, .. } = req.input; + + let objects: Vec = delete + .objects + .iter() + .map(|v| { + let version_id = v + .version_id + .as_ref() + .map(|v| match Uuid::parse_str(v) { + Ok(id) => Some(id), + Err(_) => None, + }) + .unwrap_or_default(); + ObjectToDelete { + object_name: v.key.clone(), + version_id: version_id, + } + }) + .collect(); + + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let (dobjs, _errs) = try_!(store.delete_objects(&bucket, objects, ObjectOptions::default()).await); + // info!("delete_objects res {:?} {:?}", &dobjs, errs); + + let deleted = dobjs + .iter() + .map(|v| DeletedObject { + delete_marker: { + if v.delete_marker { + Some(true) + } else { + None + } + }, + delete_marker_version_id: v.delete_marker_version_id.clone(), + key: Some(v.object_name.clone()), + version_id: v.version_id.clone(), + }) + .collect(); + + // TODO: let errors; + + let output = DeleteObjectsOutput { + deleted: Some(deleted), + // errors, + ..Default::default() + }; Ok(S3Response::new(output)) } @@ -109,7 +232,14 @@ impl S3 for FS { // mc get 1 let input = req.input; - if let Err(e) = self.store.get_bucket_info(&input.bucket, &BucketOptions {}).await { + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions {}).await { if DiskError::VolumeNotFound.is(&e) { return Err(s3_error!(NoSuchBucket)); } else { @@ -145,11 +275,14 @@ impl S3 for FS { let h = HeaderMap::new(); let opts = &ObjectOptions::default(); - let reader = try_!( - self.store - .get_object_reader(bucket.as_str(), key.as_str(), range, h, opts) - .await - ); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let reader = try_!(store.get_object_reader(bucket.as_str(), key.as_str(), range, h, opts).await); let info = reader.object_info; @@ -172,7 +305,14 @@ impl S3 for FS { async fn head_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; - if let Err(e) = self.store.get_bucket_info(&input.bucket, &BucketOptions {}).await { + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + if let Err(e) = store.get_bucket_info(&input.bucket, &BucketOptions {}).await { if DiskError::VolumeNotFound.is(&e) { return Err(s3_error!(NoSuchBucket)); } else { @@ -189,7 +329,14 @@ impl S3 for FS { // mc get 2 let HeadObjectInput { bucket, key, .. } = req.input; - let info = try_!(self.store.get_object_info(&bucket, &key, &ObjectOptions::default()).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let info = try_!(store.get_object_info(&bucket, &key, &ObjectOptions::default()).await); debug!("info {:?}", info); let content_type = try_!(ContentType::from_str("application/x-msdownload")); @@ -209,7 +356,14 @@ impl S3 for FS { async fn list_buckets(&self, _: S3Request) -> S3Result> { // mc ls - let bucket_infos = try_!(self.store.list_bucket(&BucketOptions {}).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let bucket_infos = try_!(store.list_bucket(&BucketOptions {}).await); let buckets: Vec = bucket_infos .iter() @@ -257,8 +411,15 @@ impl S3 for FS { let prefix = prefix.unwrap_or_default(); let delimiter = delimiter.unwrap_or_default(); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + let object_infos = try_!( - self.store + store .list_objects_v2( &bucket, &prefix, @@ -337,9 +498,16 @@ impl S3 for FS { let reader = PutObjReader::new(body, content_length as usize); - try_!(self.store.put_object(&bucket, &key, reader, &ObjectOptions::default()).await); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; - // self.store.put_object(bucket, object, data, opts); + try_!(store.put_object(&bucket, &key, reader, &ObjectOptions::default()).await); + + // store.put_object(bucket, object, data, opts); let output = PutObjectOutput { ..Default::default() }; Ok(S3Response::new(output)) @@ -358,11 +526,15 @@ impl S3 for FS { debug!("create_multipart_upload meta {:?}", &metadata); - let MultipartUploadResult { upload_id, .. } = try_!( - self.store - .new_multipart_upload(&bucket, &key, &ObjectOptions::default()) - .await - ); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + let MultipartUploadResult { upload_id, .. } = + try_!(store.new_multipart_upload(&bucket, &key, &ObjectOptions::default()).await); let output = CreateMultipartUploadOutput { bucket: Some(bucket), @@ -397,11 +569,14 @@ impl S3 for FS { let data = PutObjReader::new(body, content_length as usize); let opts = ObjectOptions::default(); - try_!( - self.store - .put_object_part(&bucket, &key, &upload_id, part_id, data, &opts) - .await - ); + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + + try_!(store.put_object_part(&bucket, &key, &upload_id, part_id, data, &opts).await); let output = UploadPartOutput { ..Default::default() }; Ok(S3Response::new(output)) @@ -456,8 +631,15 @@ impl S3 for FS { uploaded_parts.push(CompletePart::from(part)); } + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + try_!( - self.store + store .complete_multipart_upload(&bucket, &key, &upload_id, uploaded_parts, opts) .await ); @@ -479,9 +661,16 @@ impl S3 for FS { bucket, key, upload_id, .. } = req.input; + let layer = new_object_layer_fn(); + let lock = layer.lock().await; + let store = match lock.as_ref() { + Some(s) => s, + None => return Err(S3Error::with_message(S3ErrorCode::InternalError, format!("Not init",))), + }; + let opts = &ObjectOptions::default(); try_!( - self.store + store .abort_multipart_upload(bucket.as_str(), key.as_str(), upload_id.as_str(), opts) .await );