mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-30 16:59:52 +00:00
Merge pull request #136 from rustfs/pool-api
fix download err where versioning enable
This commit is contained in:
@@ -19,7 +19,7 @@ impl VersioningApi for VersioningConfiguration {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn prefix_enabled(&self, prefix: &str) -> bool {
|
fn prefix_enabled(&self, prefix: &str) -> bool {
|
||||||
if self.status == Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)) {
|
if self.status != Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)) {
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -399,6 +399,7 @@ impl LocalDisk {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// read xl.meta raw data
|
/// read xl.meta raw data
|
||||||
|
#[tracing::instrument(level = "debug", skip(self, volume_dir, path))]
|
||||||
async fn read_raw(
|
async fn read_raw(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -460,6 +461,7 @@ impl LocalDisk {
|
|||||||
Ok(data)
|
Ok(data)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tracing::instrument(level = "debug", skip(self, volume_dir, file_path))]
|
||||||
async fn read_all_data_with_dmtime(
|
async fn read_all_data_with_dmtime(
|
||||||
&self,
|
&self,
|
||||||
volume: &str,
|
volume: &str,
|
||||||
@@ -1851,6 +1853,7 @@ impl DiskAPI for LocalDisk {
|
|||||||
|
|
||||||
self.delete_file(&volume_dir, &xl_path, true, false).await
|
self.delete_file(&volume_dir, &xl_path, true, false).await
|
||||||
}
|
}
|
||||||
|
#[tracing::instrument(level = "debug", skip(self))]
|
||||||
async fn delete_versions(
|
async fn delete_versions(
|
||||||
&self,
|
&self,
|
||||||
volume: &str,
|
volume: &str,
|
||||||
|
|||||||
@@ -277,6 +277,7 @@ impl DiskAPI for Disk {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tracing::instrument(level = "debug", skip(self))]
|
||||||
async fn read_version(
|
async fn read_version(
|
||||||
&self,
|
&self,
|
||||||
_org_volume: &str,
|
_org_volume: &str,
|
||||||
|
|||||||
@@ -8,6 +8,52 @@ use lazy_static::lazy_static;
|
|||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
|
pub async fn del_opts(
|
||||||
|
bucket: &str,
|
||||||
|
object: &str,
|
||||||
|
vid: Option<String>,
|
||||||
|
headers: &HeaderMap<HeaderValue>,
|
||||||
|
metadata: Option<HashMap<String, String>>,
|
||||||
|
) -> Result<ObjectOptions> {
|
||||||
|
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 = put_opts_from_headers(headers, metadata)
|
||||||
|
.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.version_suspended = version_suspended;
|
||||||
|
opts.versioned = versioned;
|
||||||
|
|
||||||
|
Ok(opts)
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn put_opts(
|
pub async fn put_opts(
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
object: &str,
|
object: &str,
|
||||||
|
|||||||
@@ -273,6 +273,7 @@ struct DelObj {
|
|||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
impl ObjectIO for Sets {
|
impl ObjectIO for Sets {
|
||||||
|
#[tracing::instrument(level = "debug", skip(self))]
|
||||||
async fn get_object_reader(
|
async fn get_object_reader(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -285,6 +286,7 @@ impl ObjectIO for Sets {
|
|||||||
.get_object_reader(bucket, object, range, h, opts)
|
.get_object_reader(bucket, object, range, h, opts)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
#[tracing::instrument(level = "debug", skip(self, data))]
|
||||||
async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
||||||
self.get_disks_by_key(object).put_object(bucket, object, data, opts).await
|
self.get_disks_by_key(object).put_object(bucket, object, data, opts).await
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1083,6 +1083,7 @@ impl ObjectIO for ECStore {
|
|||||||
.get_object_reader(bucket, object.as_str(), range, h, &opts)
|
.get_object_reader(bucket, object.as_str(), range, h, &opts)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
#[tracing::instrument(level = "debug", skip(self, data))]
|
||||||
async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
||||||
check_put_object_args(bucket, object)?;
|
check_put_object_args(bucket, object)?;
|
||||||
|
|
||||||
|
|||||||
+32
-11
@@ -16,6 +16,7 @@ use ecstore::bucket::tagging::decode_tags;
|
|||||||
use ecstore::bucket::tagging::encode_tags;
|
use ecstore::bucket::tagging::encode_tags;
|
||||||
use ecstore::bucket::versioning_sys::BucketVersioningSys;
|
use ecstore::bucket::versioning_sys::BucketVersioningSys;
|
||||||
use ecstore::new_object_layer_fn;
|
use ecstore::new_object_layer_fn;
|
||||||
|
use ecstore::options::del_opts;
|
||||||
use ecstore::options::extract_metadata;
|
use ecstore::options::extract_metadata;
|
||||||
use ecstore::options::put_opts;
|
use ecstore::options::put_opts;
|
||||||
use ecstore::store_api::BucketOptions;
|
use ecstore::store_api::BucketOptions;
|
||||||
@@ -155,7 +156,14 @@ impl S3 for FS {
|
|||||||
bucket, key, version_id, ..
|
bucket, key, version_id, ..
|
||||||
} = req.input;
|
} = req.input;
|
||||||
|
|
||||||
let version_id = version_id
|
let metadata = extract_metadata(&req.headers);
|
||||||
|
|
||||||
|
let opts: ObjectOptions = del_opts(&bucket, &key, version_id, &req.headers, Some(metadata))
|
||||||
|
.await
|
||||||
|
.map_err(to_s3_error)?;
|
||||||
|
|
||||||
|
let version_id = opts
|
||||||
|
.version_id
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.map(|v| match Uuid::parse_str(v) {
|
.map(|v| match Uuid::parse_str(v) {
|
||||||
Ok(id) => Some(id),
|
Ok(id) => Some(id),
|
||||||
@@ -172,10 +180,7 @@ impl S3 for FS {
|
|||||||
let Some(store) = new_object_layer_fn() else {
|
let Some(store) = new_object_layer_fn() else {
|
||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||||
};
|
};
|
||||||
let (dobjs, _errs) = store
|
let (dobjs, _errs) = store.delete_objects(&bucket, objects, opts).await.map_err(to_s3_error)?;
|
||||||
.delete_objects(&bucket, objects, ObjectOptions::default())
|
|
||||||
.await
|
|
||||||
.map_err(to_s3_error)?;
|
|
||||||
|
|
||||||
// TODO: let errors;
|
// TODO: let errors;
|
||||||
|
|
||||||
@@ -240,11 +245,13 @@ impl S3 for FS {
|
|||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||||
};
|
};
|
||||||
|
|
||||||
let (dobjs, _errs) = store
|
let metadata = extract_metadata(&req.headers);
|
||||||
.delete_objects(&bucket, objects, ObjectOptions::default())
|
|
||||||
|
let opts: ObjectOptions = del_opts(&bucket, "", None, &req.headers, Some(metadata))
|
||||||
.await
|
.await
|
||||||
.map_err(to_s3_error)?;
|
.map_err(to_s3_error)?;
|
||||||
// info!("delete_objects res {:?} {:?}", &dobjs, errs);
|
|
||||||
|
let (dobjs, errs) = store.delete_objects(&bucket, objects, opts).await.map_err(to_s3_error)?;
|
||||||
|
|
||||||
let deleted = dobjs
|
let deleted = dobjs
|
||||||
.iter()
|
.iter()
|
||||||
@@ -263,6 +270,9 @@ impl S3 for FS {
|
|||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
// TODO: let errors;
|
// TODO: let errors;
|
||||||
|
for err in errs.iter().flatten() {
|
||||||
|
warn!("delete_objects err {:?}", err);
|
||||||
|
}
|
||||||
|
|
||||||
let output = DeleteObjectsOutput {
|
let output = DeleteObjectsOutput {
|
||||||
deleted: Some(deleted),
|
deleted: Some(deleted),
|
||||||
@@ -298,19 +308,28 @@ impl S3 for FS {
|
|||||||
async fn get_object(&self, req: S3Request<GetObjectInput>) -> S3Result<S3Response<GetObjectOutput>> {
|
async fn get_object(&self, req: S3Request<GetObjectInput>) -> S3Result<S3Response<GetObjectOutput>> {
|
||||||
// mc get 3
|
// mc get 3
|
||||||
|
|
||||||
let GetObjectInput { bucket, key, .. } = req.input;
|
let GetObjectInput {
|
||||||
|
bucket, key, version_id, ..
|
||||||
|
} = req.input;
|
||||||
|
|
||||||
let range = HTTPRangeSpec::nil();
|
let range = HTTPRangeSpec::nil();
|
||||||
|
|
||||||
let h = HeaderMap::new();
|
let h = HeaderMap::new();
|
||||||
let opts = &ObjectOptions::default();
|
|
||||||
|
let metadata = extract_metadata(&req.headers);
|
||||||
|
|
||||||
|
let opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, Some(metadata))
|
||||||
|
.await
|
||||||
|
.map_err(to_s3_error)?;
|
||||||
|
|
||||||
|
error!("get_object ObjectOptions {:?}", &opts);
|
||||||
|
|
||||||
let Some(store) = new_object_layer_fn() else {
|
let Some(store) = new_object_layer_fn() else {
|
||||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||||
};
|
};
|
||||||
|
|
||||||
let reader = store
|
let reader = store
|
||||||
.get_object_reader(bucket.as_str(), key.as_str(), range, h, opts)
|
.get_object_reader(bucket.as_str(), key.as_str(), range, h, &opts)
|
||||||
.await
|
.await
|
||||||
.map_err(to_s3_error)?;
|
.map_err(to_s3_error)?;
|
||||||
|
|
||||||
@@ -574,6 +593,8 @@ impl S3 for FS {
|
|||||||
.await
|
.await
|
||||||
.map_err(to_s3_error)?;
|
.map_err(to_s3_error)?;
|
||||||
|
|
||||||
|
error!("ObjectOptions {:?}", opts);
|
||||||
|
|
||||||
let obj_info = store
|
let obj_info = store
|
||||||
.put_object(&bucket, &key, &mut reader, &opts)
|
.put_object(&bucket, &key, &mut reader, &opts)
|
||||||
.await
|
.await
|
||||||
|
|||||||
Reference in New Issue
Block a user