From f88126008be1607028c2c87e6be38bae2fd12bfd Mon Sep 17 00:00:00 2001 From: weisd Date: Thu, 21 Nov 2024 15:59:02 +0800 Subject: [PATCH] add get_multipart_info --- ecstore/src/erasure.rs | 9 ++--- ecstore/src/sets.rs | 27 ++++++++++++-- ecstore/src/store.rs | 79 +++++++++++++++++++++++++++++++++++++++- ecstore/src/store_api.rs | 50 +++++++++++++++++-------- ecstore/src/store_err.rs | 3 ++ 5 files changed, 143 insertions(+), 25 deletions(-) diff --git a/ecstore/src/erasure.rs b/ecstore/src/erasure.rs index ee1687833..c990cce5b 100644 --- a/ecstore/src/erasure.rs +++ b/ecstore/src/erasure.rs @@ -10,7 +10,6 @@ use std::fmt::Debug; use std::io::ErrorKind; use tokio::io::DuplexStream; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tracing::debug; use tracing::warn; // use tracing::debug; use uuid::Uuid; @@ -31,10 +30,10 @@ pub struct Erasure { impl Erasure { pub fn new(data_shards: usize, parity_shards: usize, block_size: usize) -> Self { - debug!( - "Erasure new data_shards {},parity_shards {} block_size {} ", - data_shards, parity_shards, block_size - ); + // debug!( + // "Erasure new data_shards {},parity_shards {} block_size {} ", + // data_shards, parity_shards, block_size + // ); let mut encoder = None; if parity_shards > 0 { encoder = Some(ReedSolomon::new(data_shards, parity_shards).unwrap()); diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index bd8336ae8..8dd190839 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -24,8 +24,8 @@ use crate::{ set_disk::SetDisks, store_api::{ BackendInfo, BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec, - ListMultipartsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectIO, ObjectInfo, ObjectOptions, - ObjectToDelete, PartInfo, PutObjReader, StorageAPI, StorageInfo, + ListMultipartsInfo, ListObjectVersionsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartInfo, MultipartUploadResult, + ObjectIO, ObjectInfo, ObjectOptions, ObjectToDelete, PartInfo, PutObjReader, StorageAPI, StorageInfo, }, store_init::{check_format_erasure_values, get_format_erasure_in_quorum, load_format_erasure_all, save_format_file}, utils::hash, @@ -446,7 +446,17 @@ impl StorageAPI for Sets { ) -> Result { unimplemented!() } - + async fn list_object_versions( + &self, + _bucket: &str, + _prefix: &str, + _marker: &str, + _version_marker: &str, + _delimiter: &str, + _max_keys: i32, + ) -> Result { + unimplemented!() + } async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { self.get_disks_by_key(object).get_object_info(bucket, object, opts).await } @@ -513,6 +523,17 @@ impl StorageAPI for Sets { async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { self.get_disks_by_key(object).new_multipart_upload(bucket, object, opts).await } + async fn get_multipart_info( + &self, + bucket: &str, + object: &str, + upload_id: &str, + opts: &ObjectOptions, + ) -> Result { + self.get_disks_by_key(object) + .get_multipart_info(bucket, object, upload_id, opts) + .await + } async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()> { self.get_disks_by_key(object) .abort_multipart_upload(bucket, object, upload_id, opts) diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 8db66be8c..109d6ebf8 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -16,7 +16,9 @@ use crate::heal::heal_commands::{HealOpts, HealResultItem, HealScanMode, HEAL_IT use crate::heal::heal_ops::{HealEntryFn, HealSequence}; use crate::new_object_layer_fn; use crate::pools::PoolMeta; -use crate::store_api::{BackendByte, BackendDisks, BackendInfo, ListMultipartsInfo, ObjectIO, StorageInfo}; +use crate::store_api::{ + BackendByte, BackendDisks, BackendInfo, ListMultipartsInfo, ListObjectVersionsInfo, MultipartInfo, ObjectIO, StorageInfo, +}; use crate::store_err::{ is_err_bucket_exists, is_err_invalid_upload_id, is_err_object_not_found, is_err_read_quorum, is_err_version_not_found, to_object_err, StorageError, @@ -456,6 +458,11 @@ impl ECStore { ServerPoolsAvailableSpace(server_pools) } + fn is_suspended(&self, idx: usize) -> bool { + // TODO: LOCK + self.pool_meta.is_suspended(idx) + } + async fn get_pool_idx(&self, bucket: &str, object: &str, size: i64) -> Result { let idx = match self .get_pool_idx_existing_with_opts( @@ -1490,6 +1497,39 @@ impl StorageAPI for ECStore { Ok(v2) } + async fn list_object_versions( + &self, + _bucket: &str, + _prefix: &str, + marker: &str, + version_marker: &str, + _delimiter: &str, + _max_keys: i32, + ) -> Result { + if marker.is_empty() && !version_marker.is_empty() { + return Err(Error::new(StorageError::NotImplemented)); + } + + // let opts = ListPathOptions { + // bucket: bucket.to_owned(), + // marker: marker.to_owned(), + // prefix: prefix.to_owned(), + // limit: max_keys, + // ..Default::default() + // }; + + // let list = self + // .list_path(&opts, delimiter) + // .await + // .map_err(|e| to_object_err(e, vec![bucket]))?; + + // for info in list.objects.iter() { + // // + // } + + // FIXME: + unimplemented!() + } async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { check_object_args(bucket, object)?; @@ -1687,6 +1727,41 @@ impl StorageAPI for ECStore { self.pools[idx].new_multipart_upload(bucket, object, opts).await } + async fn get_multipart_info( + &self, + bucket: &str, + object: &str, + upload_id: &str, + opts: &ObjectOptions, + ) -> Result { + check_list_parts_args(bucket, object, upload_id)?; + if self.single_pool() { + return self.pools[0].get_multipart_info(bucket, object, upload_id, opts).await; + } + + for (idx, pool) in self.pools.iter().enumerate() { + if self.is_suspended(idx) { + continue; + } + + match pool.get_multipart_info(bucket, object, upload_id, opts).await { + Ok(res) => return Ok(res), + Err(err) => { + if is_err_invalid_upload_id(&err) { + continue; + } + + return Err(err); + } + } + } + + Err(Error::new(StorageError::InvalidUploadID( + bucket.to_owned(), + object.to_owned(), + upload_id.to_owned(), + ))) + } async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()> { check_abort_multipart_args(bucket, object, upload_id)?; @@ -2189,7 +2264,7 @@ fn check_put_object_part_args(bucket: &str, object: &str, upload_id: &str) -> Re check_multipart_object_args(bucket, object, upload_id) } -fn _check_list_parts_args(bucket: &str, object: &str, upload_id: &str) -> Result<()> { +fn check_list_parts_args(bucket: &str, object: &str, upload_id: &str) -> Result<()> { check_multipart_object_args(bucket, object, upload_id) } diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index eb74e8678..2853f884e 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -873,9 +873,17 @@ pub struct BackendInfo { pub drives_per_set: Vec, } +pub struct ListObjectVersionsInfo { + pub is_truncated: bool, + pub next_marker: String, + pub next_version_idmarker: String, + pub objects: Vec, + pub prefixes: Vec, +} + #[async_trait::async_trait] pub trait ObjectIO: Send + Sync + 'static { - // GetObjectNInfo + // GetObjectNInfo FIXME: async fn get_object_reader( &self, bucket: &str, @@ -890,9 +898,9 @@ pub trait ObjectIO: Send + Sync + 'static { #[async_trait::async_trait] pub trait StorageAPI: ObjectIO { - // NewNSLock - // Shutdown - // NSScanner + // NewNSLock TODO: + // Shutdown TODO: + // NSScanner TODO: async fn backend_info(&self) -> BackendInfo; async fn storage_info(&self) -> StorageInfo; @@ -902,7 +910,7 @@ pub trait StorageAPI: ObjectIO { async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result; async fn list_bucket(&self, opts: &BucketOptions) -> Result>; async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()>; - // ListObjects + // ListObjects TODO: FIXME: async fn list_objects_v2( &self, bucket: &str, @@ -913,8 +921,17 @@ pub trait StorageAPI: ObjectIO { fetch_owner: bool, start_after: &str, ) -> Result; - // ListObjectVersions - // Walk + // ListObjectVersions TODO: FIXME: + async fn list_object_versions( + &self, + bucket: &str, + prefix: &str, + marker: &str, + version_marker: &str, + delimiter: &str, + max_keys: i32, + ) -> Result; + // Walk TODO: // GetObjectNInfo ObjectIO async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result; @@ -928,8 +945,8 @@ pub trait StorageAPI: ObjectIO { opts: ObjectOptions, ) -> Result<(Vec, Vec>)>; #[warn(clippy::too_many_arguments)] - // TransitionObject - // RestoreTransitionedObject + // TransitionObject TODO: + // RestoreTransitionedObject TODO: // ListMultipartUploads async fn list_multipart_uploads( @@ -967,6 +984,13 @@ pub trait StorageAPI: ObjectIO { opts: &ObjectOptions, ) -> Result; // GetMultipartInfo + async fn get_multipart_info( + &self, + bucket: &str, + object: &str, + upload_id: &str, + opts: &ObjectOptions, + ) -> Result; // ListObjectParts async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()>; async fn complete_multipart_upload( @@ -981,12 +1005,8 @@ pub trait StorageAPI: ObjectIO { async fn get_disks(&self, pool_idx: usize, set_idx: usize) -> Result>>; // SetDriveCounts fn set_drive_counts(&self) -> Vec; - // HealFormat - // HealBucket - // HealObject - // HealObjects - // CheckAbandonedParts - // Health + + // Health TODO: // PutObjectMetadata async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result; // DecomTieredObject diff --git a/ecstore/src/store_err.rs b/ecstore/src/store_err.rs index b7f0f09ce..2ca4af46c 100644 --- a/ecstore/src/store_err.rs +++ b/ecstore/src/store_err.rs @@ -6,6 +6,9 @@ use crate::{ #[derive(Debug, thiserror::Error, PartialEq, Eq)] pub enum StorageError { + #[error("not implemented")] + NotImplemented, + #[error("Invalid arguments provided for {0}/{1}-{2}")] InvalidArgument(String, String, String),