mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-02 10:18:10 +00:00
refactor(ecstore): migrate minio/r2/rustfs warm backends to shared S3 constructor (#6776)
This commit is contained in:
@@ -19,24 +19,18 @@
|
|||||||
#![allow(clippy::all)]
|
#![allow(clippy::all)]
|
||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::Arc;
|
|
||||||
|
|
||||||
use crate::services::tier::{
|
use crate::services::tier::{
|
||||||
tier_config::TierMinIO,
|
tier_config::TierMinIO,
|
||||||
warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options},
|
warm_backend::{
|
||||||
|
S3CompatibleWarmBackendParams, TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options,
|
||||||
|
new_s3_compatible_warm_backend, optimal_part_size,
|
||||||
|
},
|
||||||
warm_backend_s3::WarmBackendS3,
|
warm_backend_s3::WarmBackendS3,
|
||||||
};
|
};
|
||||||
use rustfs_s3_client::{
|
use rustfs_s3_client::transition_api::{BucketLookupType, ReadCloser, ReaderImpl};
|
||||||
admin_handler_utils::AdminError,
|
|
||||||
api_put_object::PutObjectOptions,
|
|
||||||
credentials::{Credentials, SignatureType, Static, Value},
|
|
||||||
transition_api::{Options, ReadCloser, ReaderImpl, TransitionClient, TransitionCore},
|
|
||||||
};
|
|
||||||
use rustfs_utils::egress::validate_outbound_url;
|
use rustfs_utils::egress::validate_outbound_url;
|
||||||
use tracing::warn;
|
|
||||||
|
|
||||||
const MAX_MULTIPART_PUT_OBJECT_SIZE: i64 = 1024 * 1024 * 1024 * 1024 * 5;
|
|
||||||
const MAX_PARTS_COUNT: i64 = 10000;
|
|
||||||
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
|
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
|
||||||
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
|
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
|
||||||
|
|
||||||
@@ -44,51 +38,22 @@ pub struct WarmBackendMinIO(WarmBackendS3);
|
|||||||
|
|
||||||
impl WarmBackendMinIO {
|
impl WarmBackendMinIO {
|
||||||
pub async fn new(conf: &TierMinIO, tier: &str) -> Result<Self, std::io::Error> {
|
pub async fn new(conf: &TierMinIO, tier: &str) -> Result<Self, std::io::Error> {
|
||||||
if conf.access_key == "" || conf.secret_key == "" {
|
Ok(Self(
|
||||||
return Err(std::io::Error::other("both access and secret keys are required"));
|
new_s3_compatible_warm_backend(S3CompatibleWarmBackendParams {
|
||||||
}
|
endpoint: &conf.endpoint,
|
||||||
|
access_key: &conf.access_key,
|
||||||
if conf.bucket == "" {
|
secret_key: &conf.secret_key,
|
||||||
return Err(std::io::Error::other("no bucket name was provided"));
|
bucket: &conf.bucket,
|
||||||
}
|
prefix: &conf.prefix,
|
||||||
|
region: &conf.region,
|
||||||
let u = match url::Url::parse(&conf.endpoint) {
|
// MinIO tier endpoints are commonly path-style, so bucket addressing stays on
|
||||||
Ok(u) => u,
|
// `BucketLookupAuto`; pinning DNS here would break those deployments.
|
||||||
Err(e) => {
|
bucket_lookup: BucketLookupType::BucketLookupAuto,
|
||||||
return Err(std::io::Error::other(e.to_string()));
|
provider_tag: "minio",
|
||||||
}
|
validate_endpoint: validate_outbound_url,
|
||||||
};
|
})
|
||||||
validate_outbound_url(&u).map_err(|err| std::io::Error::other(format!("tier endpoint is not allowed: {err}")))?;
|
.await?,
|
||||||
|
))
|
||||||
let creds = Credentials::new(Static(Value {
|
|
||||||
access_key_id: conf.access_key.clone(),
|
|
||||||
secret_access_key: conf.secret_key.clone(),
|
|
||||||
session_token: "".to_string(),
|
|
||||||
signer_type: SignatureType::SignatureV4,
|
|
||||||
..Default::default()
|
|
||||||
}));
|
|
||||||
let opts = Options {
|
|
||||||
creds,
|
|
||||||
secure: u.scheme() == "https",
|
|
||||||
region: conf.region.clone(),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
let scheme = u.scheme();
|
|
||||||
let default_port = if scheme == "https" { 443 } else { 80 };
|
|
||||||
let host = u
|
|
||||||
.host_str()
|
|
||||||
.ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
|
|
||||||
let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "minio").await?;
|
|
||||||
|
|
||||||
let client = Arc::new(client);
|
|
||||||
let core = TransitionCore(Arc::clone(&client));
|
|
||||||
Ok(Self(WarmBackendS3 {
|
|
||||||
client,
|
|
||||||
core,
|
|
||||||
bucket: conf.bucket.clone(),
|
|
||||||
prefix: conf.prefix.strip_suffix("/").unwrap_or(&conf.prefix).to_owned(),
|
|
||||||
storage_class: "".to_string(),
|
|
||||||
}))
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -101,7 +66,7 @@ impl WarmBackend for WarmBackendMinIO {
|
|||||||
length: i64,
|
length: i64,
|
||||||
meta: HashMap<String, String>,
|
meta: HashMap<String, String>,
|
||||||
) -> Result<String, std::io::Error> {
|
) -> Result<String, std::io::Error> {
|
||||||
let part_size = optimal_part_size(length)?;
|
let part_size = optimal_part_size(length, MIN_PART_SIZE)?;
|
||||||
let client = self.0.client.clone();
|
let client = self.0.client.clone();
|
||||||
let res = client
|
let res = client
|
||||||
.put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &{
|
.put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &{
|
||||||
@@ -150,32 +115,15 @@ impl crate::services::tier::warm_backend::TransitionCandidateReconciler for Warm
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn optimal_part_size(object_size: i64) -> Result<i64, std::io::Error> {
|
|
||||||
let mut object_size = object_size;
|
|
||||||
if object_size == -1 {
|
|
||||||
object_size = MAX_MULTIPART_PUT_OBJECT_SIZE;
|
|
||||||
}
|
|
||||||
|
|
||||||
if object_size > MAX_MULTIPART_PUT_OBJECT_SIZE {
|
|
||||||
return Err(std::io::Error::other("entity too large"));
|
|
||||||
}
|
|
||||||
|
|
||||||
let configured_part_size = MIN_PART_SIZE;
|
|
||||||
let mut part_size_flt = object_size as f64 / MAX_PARTS_COUNT as f64;
|
|
||||||
part_size_flt = (part_size_flt as f64 / configured_part_size as f64).ceil() * configured_part_size as f64;
|
|
||||||
|
|
||||||
let part_size = part_size_flt as i64;
|
|
||||||
if part_size == 0 {
|
|
||||||
return Ok(MIN_PART_SIZE);
|
|
||||||
}
|
|
||||||
Ok(part_size)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::services::tier::tier_config::TierMinIO;
|
use crate::services::tier::tier_config::TierMinIO;
|
||||||
|
|
||||||
|
/// The SSRF guard itself is exercised once, generically, in
|
||||||
|
/// `warm_backend::tests` (see backlog#2040/backlog#2043 and
|
||||||
|
/// rustfs/rustfs#6764) — this test only pins that this provider's
|
||||||
|
/// production constructor really is wired through that shared path.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn new_rejects_loopback_endpoint_before_network_setup() {
|
async fn new_rejects_loopback_endpoint_before_network_setup() {
|
||||||
let conf = TierMinIO {
|
let conf = TierMinIO {
|
||||||
|
|||||||
@@ -19,24 +19,18 @@
|
|||||||
#![allow(clippy::all)]
|
#![allow(clippy::all)]
|
||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::Arc;
|
|
||||||
|
|
||||||
use crate::services::tier::{
|
use crate::services::tier::{
|
||||||
tier_config::TierR2,
|
tier_config::TierR2,
|
||||||
warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options},
|
warm_backend::{
|
||||||
|
S3CompatibleWarmBackendParams, TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options,
|
||||||
|
new_s3_compatible_warm_backend, optimal_part_size,
|
||||||
|
},
|
||||||
warm_backend_s3::WarmBackendS3,
|
warm_backend_s3::WarmBackendS3,
|
||||||
};
|
};
|
||||||
use rustfs_s3_client::{
|
use rustfs_s3_client::transition_api::{BucketLookupType, ReadCloser, ReaderImpl};
|
||||||
admin_handler_utils::AdminError,
|
|
||||||
api_put_object::PutObjectOptions,
|
|
||||||
credentials::{Credentials, SignatureType, Static, Value},
|
|
||||||
transition_api::{Options, ReadCloser, ReaderImpl, TransitionClient, TransitionCore},
|
|
||||||
};
|
|
||||||
use rustfs_utils::egress::validate_outbound_url;
|
use rustfs_utils::egress::validate_outbound_url;
|
||||||
use tracing::warn;
|
|
||||||
|
|
||||||
const MAX_MULTIPART_PUT_OBJECT_SIZE: i64 = 1024 * 1024 * 1024 * 1024 * 5;
|
|
||||||
const MAX_PARTS_COUNT: i64 = 10000;
|
|
||||||
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
|
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
|
||||||
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
|
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
|
||||||
|
|
||||||
@@ -44,51 +38,22 @@ pub struct WarmBackendR2(WarmBackendS3);
|
|||||||
|
|
||||||
impl WarmBackendR2 {
|
impl WarmBackendR2 {
|
||||||
pub async fn new(conf: &TierR2, tier: &str) -> Result<Self, std::io::Error> {
|
pub async fn new(conf: &TierR2, tier: &str) -> Result<Self, std::io::Error> {
|
||||||
if conf.access_key == "" || conf.secret_key == "" {
|
Ok(Self(
|
||||||
return Err(std::io::Error::other("both access and secret keys are required"));
|
new_s3_compatible_warm_backend(S3CompatibleWarmBackendParams {
|
||||||
}
|
endpoint: &conf.endpoint,
|
||||||
|
access_key: &conf.access_key,
|
||||||
if conf.bucket == "" {
|
secret_key: &conf.secret_key,
|
||||||
return Err(std::io::Error::other("no bucket name was provided"));
|
bucket: &conf.bucket,
|
||||||
}
|
prefix: &conf.prefix,
|
||||||
|
region: &conf.region,
|
||||||
let u = match url::Url::parse(&conf.endpoint) {
|
// R2 tier endpoints are commonly path-style, so bucket addressing stays on
|
||||||
Ok(u) => u,
|
// `BucketLookupAuto`; pinning DNS here would break those deployments.
|
||||||
Err(e) => {
|
bucket_lookup: BucketLookupType::BucketLookupAuto,
|
||||||
return Err(std::io::Error::other(e.to_string()));
|
provider_tag: "r2",
|
||||||
}
|
validate_endpoint: validate_outbound_url,
|
||||||
};
|
})
|
||||||
validate_outbound_url(&u).map_err(|err| std::io::Error::other(format!("tier endpoint is not allowed: {err}")))?;
|
.await?,
|
||||||
|
))
|
||||||
let creds = Credentials::new(Static(Value {
|
|
||||||
access_key_id: conf.access_key.clone(),
|
|
||||||
secret_access_key: conf.secret_key.clone(),
|
|
||||||
session_token: "".to_string(),
|
|
||||||
signer_type: SignatureType::SignatureV4,
|
|
||||||
..Default::default()
|
|
||||||
}));
|
|
||||||
let opts = Options {
|
|
||||||
creds,
|
|
||||||
secure: u.scheme() == "https",
|
|
||||||
region: conf.region.clone(),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
let scheme = u.scheme();
|
|
||||||
let default_port = if scheme == "https" { 443 } else { 80 };
|
|
||||||
let host = u
|
|
||||||
.host_str()
|
|
||||||
.ok_or_else(|| std::io::Error::other("Invalid endpoint URL: missing host"))?;
|
|
||||||
let client = TransitionClient::new(&format!("{}:{}", host, u.port().unwrap_or(default_port)), opts, "r2").await?;
|
|
||||||
|
|
||||||
let client = Arc::new(client);
|
|
||||||
let core = TransitionCore(Arc::clone(&client));
|
|
||||||
Ok(Self(WarmBackendS3 {
|
|
||||||
client,
|
|
||||||
core,
|
|
||||||
bucket: conf.bucket.clone(),
|
|
||||||
prefix: conf.prefix.strip_suffix("/").unwrap_or(&conf.prefix).to_owned(),
|
|
||||||
storage_class: "".to_string(),
|
|
||||||
}))
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -101,7 +66,7 @@ impl WarmBackend for WarmBackendR2 {
|
|||||||
length: i64,
|
length: i64,
|
||||||
meta: HashMap<String, String>,
|
meta: HashMap<String, String>,
|
||||||
) -> Result<String, std::io::Error> {
|
) -> Result<String, std::io::Error> {
|
||||||
let part_size = optimal_part_size(length)?;
|
let part_size = optimal_part_size(length, MIN_PART_SIZE)?;
|
||||||
let client = self.0.client.clone();
|
let client = self.0.client.clone();
|
||||||
let res = client
|
let res = client
|
||||||
.put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &{
|
.put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &{
|
||||||
@@ -150,32 +115,15 @@ impl crate::services::tier::warm_backend::TransitionCandidateReconciler for Warm
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn optimal_part_size(object_size: i64) -> Result<i64, std::io::Error> {
|
|
||||||
let mut object_size = object_size;
|
|
||||||
if object_size == -1 {
|
|
||||||
object_size = MAX_MULTIPART_PUT_OBJECT_SIZE;
|
|
||||||
}
|
|
||||||
|
|
||||||
if object_size > MAX_MULTIPART_PUT_OBJECT_SIZE {
|
|
||||||
return Err(std::io::Error::other("entity too large"));
|
|
||||||
}
|
|
||||||
|
|
||||||
let configured_part_size = MIN_PART_SIZE;
|
|
||||||
let mut part_size_flt = object_size as f64 / MAX_PARTS_COUNT as f64;
|
|
||||||
part_size_flt = (part_size_flt as f64 / configured_part_size as f64).ceil() * configured_part_size as f64;
|
|
||||||
|
|
||||||
let part_size = part_size_flt as i64;
|
|
||||||
if part_size == 0 {
|
|
||||||
return Ok(MIN_PART_SIZE);
|
|
||||||
}
|
|
||||||
Ok(part_size)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::services::tier::tier_config::TierR2;
|
use crate::services::tier::tier_config::TierR2;
|
||||||
|
|
||||||
|
/// The SSRF guard itself is exercised once, generically, in
|
||||||
|
/// `warm_backend::tests` (see backlog#2040/backlog#2043 and
|
||||||
|
/// rustfs/rustfs#6764) — this test only pins that this provider's
|
||||||
|
/// production constructor really is wired through that shared path.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn new_rejects_loopback_endpoint_before_network_setup() {
|
async fn new_rejects_loopback_endpoint_before_network_setup() {
|
||||||
let conf = TierR2 {
|
let conf = TierR2 {
|
||||||
|
|||||||
@@ -19,23 +19,18 @@
|
|||||||
#![allow(clippy::all)]
|
#![allow(clippy::all)]
|
||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::Arc;
|
|
||||||
|
|
||||||
use crate::services::tier::{
|
use crate::services::tier::{
|
||||||
tier_config::TierRustFS,
|
tier_config::TierRustFS,
|
||||||
warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options},
|
warm_backend::{
|
||||||
|
S3CompatibleWarmBackendParams, TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options,
|
||||||
|
new_s3_compatible_warm_backend, optimal_part_size,
|
||||||
|
},
|
||||||
warm_backend_s3::WarmBackendS3,
|
warm_backend_s3::WarmBackendS3,
|
||||||
};
|
};
|
||||||
use rustfs_s3_client::{
|
use rustfs_s3_client::transition_api::{BucketLookupType, ReadCloser, ReaderImpl};
|
||||||
admin_handler_utils::AdminError,
|
|
||||||
api_put_object::PutObjectOptions,
|
|
||||||
credentials::{Credentials, SignatureType, Static, Value},
|
|
||||||
transition_api::{Options, ReadCloser, ReaderImpl, TransitionClient, TransitionCore},
|
|
||||||
};
|
|
||||||
use rustfs_utils::egress::{OutboundUrlError, validate_outbound_url};
|
use rustfs_utils::egress::{OutboundUrlError, validate_outbound_url};
|
||||||
|
|
||||||
const MAX_MULTIPART_PUT_OBJECT_SIZE: i64 = 1024 * 1024 * 1024 * 1024 * 5;
|
|
||||||
const MAX_PARTS_COUNT: i64 = 10000;
|
|
||||||
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
|
const _MAX_PART_SIZE: i64 = 1024 * 1024 * 1024 * 5;
|
||||||
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
|
const MIN_PART_SIZE: i64 = 1024 * 1024 * 128;
|
||||||
// Debug-only opt-in for single-host test/dev setups; release builds always reject loopback.
|
// Debug-only opt-in for single-host test/dev setups; release builds always reject loopback.
|
||||||
@@ -63,11 +58,16 @@ pub struct WarmBackendRustFS(WarmBackendS3);
|
|||||||
|
|
||||||
impl WarmBackendRustFS {
|
impl WarmBackendRustFS {
|
||||||
pub async fn new(conf: &TierRustFS, tier: &str) -> Result<Self, std::io::Error> {
|
pub async fn new(conf: &TierRustFS, tier: &str) -> Result<Self, std::io::Error> {
|
||||||
if conf.access_key == "" || conf.secret_key == "" {
|
// This provider reports endpoint problems with its own wording (and keeps the
|
||||||
|
// `url::ParseError` as the io::Error source) while the shared constructor carries the
|
||||||
|
// MinIO-derived texts. Unifying the two is a separate change, so the endpoint is
|
||||||
|
// pre-validated here, after the credential and bucket checks so the order in which the
|
||||||
|
// shared constructor would report the same failures is preserved.
|
||||||
|
if conf.access_key.is_empty() || conf.secret_key.is_empty() {
|
||||||
return Err(std::io::Error::other("both access and secret keys are required"));
|
return Err(std::io::Error::other("both access and secret keys are required"));
|
||||||
}
|
}
|
||||||
|
|
||||||
if conf.bucket == "" {
|
if conf.bucket.is_empty() {
|
||||||
return Err(std::io::Error::other("no bucket name was provided"));
|
return Err(std::io::Error::other("no bucket name was provided"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -75,37 +75,30 @@ impl WarmBackendRustFS {
|
|||||||
Ok(u) => u,
|
Ok(u) => u,
|
||||||
Err(e) => return Err(std::io::Error::other(e)),
|
Err(e) => return Err(std::io::Error::other(e)),
|
||||||
};
|
};
|
||||||
validate_rustfs_tier_endpoint(&u).map_err(|err| std::io::Error::other(format!("tier endpoint is not allowed: {err}")))?;
|
|
||||||
|
|
||||||
let creds = Credentials::new(Static(Value {
|
if u.host_str().is_none() {
|
||||||
access_key_id: conf.access_key.clone(),
|
return Err(std::io::Error::other("endpoint URL must include a host"));
|
||||||
secret_access_key: conf.secret_key.clone(),
|
}
|
||||||
session_token: "".to_string(),
|
|
||||||
signer_type: SignatureType::SignatureV4,
|
|
||||||
..Default::default()
|
|
||||||
}));
|
|
||||||
let opts = Options {
|
|
||||||
creds,
|
|
||||||
secure: u.scheme() == "https",
|
|
||||||
region: conf.region.clone(),
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
let scheme = u.scheme();
|
|
||||||
let default_port = if scheme == "https" { 443 } else { 80 };
|
|
||||||
let host = u
|
|
||||||
.host_str()
|
|
||||||
.ok_or_else(|| std::io::Error::other("endpoint URL must include a host"))?;
|
|
||||||
let client = TransitionClient::new(&format!("{host}:{}", u.port().unwrap_or(default_port)), opts, "rustfs").await?;
|
|
||||||
|
|
||||||
let client = Arc::new(client);
|
Ok(Self(
|
||||||
let core = TransitionCore(Arc::clone(&client));
|
new_s3_compatible_warm_backend(S3CompatibleWarmBackendParams {
|
||||||
Ok(Self(WarmBackendS3 {
|
endpoint: &conf.endpoint,
|
||||||
client,
|
access_key: &conf.access_key,
|
||||||
core,
|
secret_key: &conf.secret_key,
|
||||||
bucket: conf.bucket.clone(),
|
bucket: &conf.bucket,
|
||||||
prefix: conf.prefix.strip_suffix("/").unwrap_or(&conf.prefix).to_owned(),
|
prefix: &conf.prefix,
|
||||||
storage_class: "".to_string(),
|
region: &conf.region,
|
||||||
}))
|
// RustFS tier endpoints are path-style, so bucket addressing stays on
|
||||||
|
// `BucketLookupAuto`; pinning DNS here would break those endpoints.
|
||||||
|
bucket_lookup: BucketLookupType::BucketLookupAuto,
|
||||||
|
provider_tag: "rustfs",
|
||||||
|
// Debug-only, env-gated loopback exception for this provider's own e2e tier
|
||||||
|
// tests (rustfs/rustfs#6773); every other provider passes plain
|
||||||
|
// `validate_outbound_url`.
|
||||||
|
validate_endpoint: validate_rustfs_tier_endpoint,
|
||||||
|
})
|
||||||
|
.await?,
|
||||||
|
))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -118,7 +111,7 @@ impl WarmBackend for WarmBackendRustFS {
|
|||||||
length: i64,
|
length: i64,
|
||||||
meta: HashMap<String, String>,
|
meta: HashMap<String, String>,
|
||||||
) -> Result<String, std::io::Error> {
|
) -> Result<String, std::io::Error> {
|
||||||
let part_size = optimal_part_size(length)?;
|
let part_size = optimal_part_size(length, MIN_PART_SIZE)?;
|
||||||
let client = self.0.client.clone();
|
let client = self.0.client.clone();
|
||||||
let res = client
|
let res = client
|
||||||
.put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &{
|
.put_object(&self.0.bucket, &self.0.get_dest(object), r, length, &{
|
||||||
@@ -167,27 +160,6 @@ impl crate::services::tier::warm_backend::TransitionCandidateReconciler for Warm
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn optimal_part_size(object_size: i64) -> Result<i64, std::io::Error> {
|
|
||||||
let mut object_size = object_size;
|
|
||||||
if object_size == -1 {
|
|
||||||
object_size = MAX_MULTIPART_PUT_OBJECT_SIZE;
|
|
||||||
}
|
|
||||||
|
|
||||||
if object_size > MAX_MULTIPART_PUT_OBJECT_SIZE {
|
|
||||||
return Err(std::io::Error::other("entity too large"));
|
|
||||||
}
|
|
||||||
|
|
||||||
let configured_part_size = MIN_PART_SIZE;
|
|
||||||
let mut part_size_flt = object_size as f64 / MAX_PARTS_COUNT as f64;
|
|
||||||
part_size_flt = (part_size_flt as f64 / configured_part_size as f64).ceil() * configured_part_size as f64;
|
|
||||||
|
|
||||||
let part_size = part_size_flt as i64;
|
|
||||||
if part_size == 0 {
|
|
||||||
return Ok(MIN_PART_SIZE);
|
|
||||||
}
|
|
||||||
Ok(part_size)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use futures::FutureExt;
|
use futures::FutureExt;
|
||||||
|
|||||||
@@ -56,9 +56,6 @@
|
|||||||
1|crates/ecstore/src/services/tier/tier_config.rs
|
1|crates/ecstore/src/services/tier/tier_config.rs
|
||||||
1|crates/ecstore/src/services/tier/warm_backend.rs
|
1|crates/ecstore/src/services/tier/warm_backend.rs
|
||||||
2|crates/ecstore/src/services/tier/warm_backend_gcs.rs
|
2|crates/ecstore/src/services/tier/warm_backend_gcs.rs
|
||||||
1|crates/ecstore/src/services/tier/warm_backend_minio.rs
|
|
||||||
1|crates/ecstore/src/services/tier/warm_backend_r2.rs
|
|
||||||
1|crates/ecstore/src/services/tier/warm_backend_rustfs.rs
|
|
||||||
1|crates/ecstore/src/services/tier/warm_backend_s3.rs
|
1|crates/ecstore/src/services/tier/warm_backend_s3.rs
|
||||||
1|crates/ecstore/src/services/tier/warm_backend_wasabi.rs
|
1|crates/ecstore/src/services/tier/warm_backend_wasabi.rs
|
||||||
7|crates/ecstore/src/set_disk/core/io_primitives.rs
|
7|crates/ecstore/src/set_disk/core/io_primitives.rs
|
||||||
|
|||||||
Reference in New Issue
Block a user