mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
add get_multipart_info
This commit is contained in:
@@ -10,7 +10,6 @@ use std::fmt::Debug;
|
|||||||
use std::io::ErrorKind;
|
use std::io::ErrorKind;
|
||||||
use tokio::io::DuplexStream;
|
use tokio::io::DuplexStream;
|
||||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
use tracing::debug;
|
|
||||||
use tracing::warn;
|
use tracing::warn;
|
||||||
// use tracing::debug;
|
// use tracing::debug;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
@@ -31,10 +30,10 @@ pub struct Erasure {
|
|||||||
|
|
||||||
impl Erasure {
|
impl Erasure {
|
||||||
pub fn new(data_shards: usize, parity_shards: usize, block_size: usize) -> Self {
|
pub fn new(data_shards: usize, parity_shards: usize, block_size: usize) -> Self {
|
||||||
debug!(
|
// debug!(
|
||||||
"Erasure new data_shards {},parity_shards {} block_size {} ",
|
// "Erasure new data_shards {},parity_shards {} block_size {} ",
|
||||||
data_shards, parity_shards, block_size
|
// data_shards, parity_shards, block_size
|
||||||
);
|
// );
|
||||||
let mut encoder = None;
|
let mut encoder = None;
|
||||||
if parity_shards > 0 {
|
if parity_shards > 0 {
|
||||||
encoder = Some(ReedSolomon::new(data_shards, parity_shards).unwrap());
|
encoder = Some(ReedSolomon::new(data_shards, parity_shards).unwrap());
|
||||||
|
|||||||
+24
-3
@@ -24,8 +24,8 @@ use crate::{
|
|||||||
set_disk::SetDisks,
|
set_disk::SetDisks,
|
||||||
store_api::{
|
store_api::{
|
||||||
BackendInfo, BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec,
|
BackendInfo, BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec,
|
||||||
ListMultipartsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectIO, ObjectInfo, ObjectOptions,
|
ListMultipartsInfo, ListObjectVersionsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartInfo, MultipartUploadResult,
|
||||||
ObjectToDelete, PartInfo, PutObjReader, StorageAPI, StorageInfo,
|
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},
|
store_init::{check_format_erasure_values, get_format_erasure_in_quorum, load_format_erasure_all, save_format_file},
|
||||||
utils::hash,
|
utils::hash,
|
||||||
@@ -446,7 +446,17 @@ impl StorageAPI for Sets {
|
|||||||
) -> Result<ListObjectsV2Info> {
|
) -> Result<ListObjectsV2Info> {
|
||||||
unimplemented!()
|
unimplemented!()
|
||||||
}
|
}
|
||||||
|
async fn list_object_versions(
|
||||||
|
&self,
|
||||||
|
_bucket: &str,
|
||||||
|
_prefix: &str,
|
||||||
|
_marker: &str,
|
||||||
|
_version_marker: &str,
|
||||||
|
_delimiter: &str,
|
||||||
|
_max_keys: i32,
|
||||||
|
) -> Result<ListObjectVersionsInfo> {
|
||||||
|
unimplemented!()
|
||||||
|
}
|
||||||
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
||||||
self.get_disks_by_key(object).get_object_info(bucket, object, opts).await
|
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<MultipartUploadResult> {
|
async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<MultipartUploadResult> {
|
||||||
self.get_disks_by_key(object).new_multipart_upload(bucket, object, opts).await
|
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<MultipartInfo> {
|
||||||
|
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<()> {
|
async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()> {
|
||||||
self.get_disks_by_key(object)
|
self.get_disks_by_key(object)
|
||||||
.abort_multipart_upload(bucket, object, upload_id, opts)
|
.abort_multipart_upload(bucket, object, upload_id, opts)
|
||||||
|
|||||||
+77
-2
@@ -16,7 +16,9 @@ use crate::heal::heal_commands::{HealOpts, HealResultItem, HealScanMode, HEAL_IT
|
|||||||
use crate::heal::heal_ops::{HealEntryFn, HealSequence};
|
use crate::heal::heal_ops::{HealEntryFn, HealSequence};
|
||||||
use crate::new_object_layer_fn;
|
use crate::new_object_layer_fn;
|
||||||
use crate::pools::PoolMeta;
|
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::{
|
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,
|
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,
|
to_object_err, StorageError,
|
||||||
@@ -456,6 +458,11 @@ impl ECStore {
|
|||||||
ServerPoolsAvailableSpace(server_pools)
|
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<usize> {
|
async fn get_pool_idx(&self, bucket: &str, object: &str, size: i64) -> Result<usize> {
|
||||||
let idx = match self
|
let idx = match self
|
||||||
.get_pool_idx_existing_with_opts(
|
.get_pool_idx_existing_with_opts(
|
||||||
@@ -1490,6 +1497,39 @@ impl StorageAPI for ECStore {
|
|||||||
|
|
||||||
Ok(v2)
|
Ok(v2)
|
||||||
}
|
}
|
||||||
|
async fn list_object_versions(
|
||||||
|
&self,
|
||||||
|
_bucket: &str,
|
||||||
|
_prefix: &str,
|
||||||
|
marker: &str,
|
||||||
|
version_marker: &str,
|
||||||
|
_delimiter: &str,
|
||||||
|
_max_keys: i32,
|
||||||
|
) -> Result<ListObjectVersionsInfo> {
|
||||||
|
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<ObjectInfo> {
|
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
||||||
check_object_args(bucket, object)?;
|
check_object_args(bucket, object)?;
|
||||||
|
|
||||||
@@ -1687,6 +1727,41 @@ impl StorageAPI for ECStore {
|
|||||||
|
|
||||||
self.pools[idx].new_multipart_upload(bucket, object, opts).await
|
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<MultipartInfo> {
|
||||||
|
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<()> {
|
async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()> {
|
||||||
check_abort_multipart_args(bucket, object, upload_id)?;
|
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)
|
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)
|
check_multipart_object_args(bucket, object, upload_id)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+35
-15
@@ -873,9 +873,17 @@ pub struct BackendInfo {
|
|||||||
pub drives_per_set: Vec<usize>,
|
pub drives_per_set: Vec<usize>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub struct ListObjectVersionsInfo {
|
||||||
|
pub is_truncated: bool,
|
||||||
|
pub next_marker: String,
|
||||||
|
pub next_version_idmarker: String,
|
||||||
|
pub objects: Vec<ObjectInfo>,
|
||||||
|
pub prefixes: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
pub trait ObjectIO: Send + Sync + 'static {
|
pub trait ObjectIO: Send + Sync + 'static {
|
||||||
// GetObjectNInfo
|
// GetObjectNInfo FIXME:
|
||||||
async fn get_object_reader(
|
async fn get_object_reader(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -890,9 +898,9 @@ pub trait ObjectIO: Send + Sync + 'static {
|
|||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
pub trait StorageAPI: ObjectIO {
|
pub trait StorageAPI: ObjectIO {
|
||||||
// NewNSLock
|
// NewNSLock TODO:
|
||||||
// Shutdown
|
// Shutdown TODO:
|
||||||
// NSScanner
|
// NSScanner TODO:
|
||||||
|
|
||||||
async fn backend_info(&self) -> BackendInfo;
|
async fn backend_info(&self) -> BackendInfo;
|
||||||
async fn storage_info(&self) -> StorageInfo;
|
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<BucketInfo>;
|
async fn get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo>;
|
||||||
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>>;
|
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>>;
|
||||||
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()>;
|
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()>;
|
||||||
// ListObjects
|
// ListObjects TODO: FIXME:
|
||||||
async fn list_objects_v2(
|
async fn list_objects_v2(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -913,8 +921,17 @@ pub trait StorageAPI: ObjectIO {
|
|||||||
fetch_owner: bool,
|
fetch_owner: bool,
|
||||||
start_after: &str,
|
start_after: &str,
|
||||||
) -> Result<ListObjectsV2Info>;
|
) -> Result<ListObjectsV2Info>;
|
||||||
// ListObjectVersions
|
// ListObjectVersions TODO: FIXME:
|
||||||
// Walk
|
async fn list_object_versions(
|
||||||
|
&self,
|
||||||
|
bucket: &str,
|
||||||
|
prefix: &str,
|
||||||
|
marker: &str,
|
||||||
|
version_marker: &str,
|
||||||
|
delimiter: &str,
|
||||||
|
max_keys: i32,
|
||||||
|
) -> Result<ListObjectVersionsInfo>;
|
||||||
|
// Walk TODO:
|
||||||
|
|
||||||
// GetObjectNInfo ObjectIO
|
// GetObjectNInfo ObjectIO
|
||||||
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo>;
|
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo>;
|
||||||
@@ -928,8 +945,8 @@ pub trait StorageAPI: ObjectIO {
|
|||||||
opts: ObjectOptions,
|
opts: ObjectOptions,
|
||||||
) -> Result<(Vec<DeletedObject>, Vec<Option<Error>>)>;
|
) -> Result<(Vec<DeletedObject>, Vec<Option<Error>>)>;
|
||||||
#[warn(clippy::too_many_arguments)]
|
#[warn(clippy::too_many_arguments)]
|
||||||
// TransitionObject
|
// TransitionObject TODO:
|
||||||
// RestoreTransitionedObject
|
// RestoreTransitionedObject TODO:
|
||||||
|
|
||||||
// ListMultipartUploads
|
// ListMultipartUploads
|
||||||
async fn list_multipart_uploads(
|
async fn list_multipart_uploads(
|
||||||
@@ -967,6 +984,13 @@ pub trait StorageAPI: ObjectIO {
|
|||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
) -> Result<PartInfo>;
|
) -> Result<PartInfo>;
|
||||||
// GetMultipartInfo
|
// GetMultipartInfo
|
||||||
|
async fn get_multipart_info(
|
||||||
|
&self,
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
upload_id: &str,
|
||||||
|
opts: &ObjectOptions,
|
||||||
|
) -> Result<MultipartInfo>;
|
||||||
// ListObjectParts
|
// ListObjectParts
|
||||||
async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()>;
|
async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()>;
|
||||||
async fn complete_multipart_upload(
|
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<Vec<Option<DiskStore>>>;
|
async fn get_disks(&self, pool_idx: usize, set_idx: usize) -> Result<Vec<Option<DiskStore>>>;
|
||||||
// SetDriveCounts
|
// SetDriveCounts
|
||||||
fn set_drive_counts(&self) -> Vec<usize>;
|
fn set_drive_counts(&self) -> Vec<usize>;
|
||||||
// HealFormat
|
|
||||||
// HealBucket
|
// Health TODO:
|
||||||
// HealObject
|
|
||||||
// HealObjects
|
|
||||||
// CheckAbandonedParts
|
|
||||||
// Health
|
|
||||||
// PutObjectMetadata
|
// PutObjectMetadata
|
||||||
async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo>;
|
async fn put_object_metadata(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo>;
|
||||||
// DecomTieredObject
|
// DecomTieredObject
|
||||||
|
|||||||
@@ -6,6 +6,9 @@ use crate::{
|
|||||||
|
|
||||||
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
|
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
|
||||||
pub enum StorageError {
|
pub enum StorageError {
|
||||||
|
#[error("not implemented")]
|
||||||
|
NotImplemented,
|
||||||
|
|
||||||
#[error("Invalid arguments provided for {0}/{1}-{2}")]
|
#[error("Invalid arguments provided for {0}/{1}-{2}")]
|
||||||
InvalidArgument(String, String, String),
|
InvalidArgument(String, String, String),
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user