mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-01 17:58:22 +00:00
fix(storage): gate multipart upload part pressure (#6781)
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -297,7 +297,7 @@ const _: () = assert!(!DEFAULT_PUT_FOREGROUND_ADMISSION_ENABLE);
|
|||||||
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE";
|
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE";
|
||||||
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: bool = true;
|
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: bool = true;
|
||||||
|
|
||||||
/// Maximum large foreground PutObject requests admitted concurrently per process.
|
/// Maximum automatic foreground write requests admitted concurrently per process.
|
||||||
///
|
///
|
||||||
/// `0` derives a conservative default from the local disk-read scheduler cap,
|
/// `0` derives a conservative default from the local disk-read scheduler cap,
|
||||||
/// currently clamped to protect the commit path without making ordinary high
|
/// currently clamped to protect the commit path without making ordinary high
|
||||||
@@ -305,14 +305,24 @@ pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_ENABLE: bool = true;
|
|||||||
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT";
|
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT";
|
||||||
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT: usize = 0;
|
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_LIMIT: usize = 0;
|
||||||
|
|
||||||
/// Minimum object size that enters automatic large PutObject admission.
|
/// Minimum direct PutObject size that enters automatic foreground write admission.
|
||||||
///
|
///
|
||||||
/// Requests with an unknown size are treated as large because the write pressure
|
/// Requests with an unknown size are treated as large because the write pressure
|
||||||
/// cannot be bounded from headers.
|
/// cannot be bounded from headers.
|
||||||
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES";
|
pub const ENV_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES: &str = "RUSTFS_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES";
|
||||||
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES: usize = 32 * 1024 * 1024;
|
pub const DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES: usize = 32 * 1024 * 1024;
|
||||||
|
|
||||||
/// Time in milliseconds a large foreground PutObject waits for a permit.
|
/// Minimum UploadPart size that enters automatic foreground write admission.
|
||||||
|
///
|
||||||
|
/// Multipart pressure is often many moderate-sized parts rather than one very
|
||||||
|
/// large request. The default gates every multipart part through the same permit
|
||||||
|
/// pool as large/unknown-size PutObject while keeping small direct PUTs on the
|
||||||
|
/// legacy path.
|
||||||
|
pub const ENV_PUT_MULTIPART_FOREGROUND_ADMISSION_MIN_SIZE_BYTES: &str =
|
||||||
|
"RUSTFS_PUT_MULTIPART_FOREGROUND_ADMISSION_MIN_SIZE_BYTES";
|
||||||
|
pub const DEFAULT_PUT_MULTIPART_FOREGROUND_ADMISSION_MIN_SIZE_BYTES: usize = 0;
|
||||||
|
|
||||||
|
/// Time in milliseconds an automatic foreground write waits for a permit.
|
||||||
///
|
///
|
||||||
/// A short wait smooths transient bursts while still returning S3
|
/// A short wait smooths transient bursts while still returning S3
|
||||||
/// `SlowDown`/503 before body ingest when the node is already saturated.
|
/// `SlowDown`/503 before body ingest when the node is already saturated.
|
||||||
|
|||||||
@@ -67,7 +67,10 @@ use super::storage_api::multipart_usecase::sse::{
|
|||||||
use super::storage_api::multipart_usecase::{
|
use super::storage_api::multipart_usecase::{
|
||||||
StorageObjectInfo as ObjectInfo, StorageObjectOptions as ObjectOptions, StoragePutObjReader as PutObjReader,
|
StorageObjectInfo as ObjectInfo, StorageObjectOptions as ObjectOptions, StoragePutObjReader as PutObjReader,
|
||||||
};
|
};
|
||||||
use crate::app::object::{guard_put_object_body_read_timeout, put_object_body_read_timeout};
|
use crate::app::object::{
|
||||||
|
ConcurrencyManager, ForegroundWriteAdmission, get_concurrency_manager, guard_put_object_body_read_timeout,
|
||||||
|
put_object_body_read_timeout,
|
||||||
|
};
|
||||||
use crate::app::object_data_cache::{
|
use crate::app::object_data_cache::{
|
||||||
ObjectDataCacheAdapter, invalidate_object_data_cache_after_complete_multipart_success,
|
ObjectDataCacheAdapter, invalidate_object_data_cache_after_complete_multipart_success,
|
||||||
invalidate_object_data_cache_before_mutation,
|
invalidate_object_data_cache_before_mutation,
|
||||||
@@ -90,6 +93,7 @@ use crate::table_catalog;
|
|||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use futures::StreamExt;
|
use futures::StreamExt;
|
||||||
use http::{HeaderMap, HeaderValue, Uri};
|
use http::{HeaderMap, HeaderValue, Uri};
|
||||||
|
use metrics::counter;
|
||||||
use rustfs_io_metrics::record_s3_op;
|
use rustfs_io_metrics::record_s3_op;
|
||||||
use rustfs_s3_ops::S3Operation;
|
use rustfs_s3_ops::S3Operation;
|
||||||
use rustfs_targets::EventName;
|
use rustfs_targets::EventName;
|
||||||
@@ -382,17 +386,24 @@ fn build_complete_multipart_location(headers: &HeaderMap, uri: &Uri, bucket: &st
|
|||||||
#[derive(Clone, Default)]
|
#[derive(Clone, Default)]
|
||||||
pub struct DefaultMultipartUsecase {
|
pub struct DefaultMultipartUsecase {
|
||||||
context: Option<Arc<AppContext>>,
|
context: Option<Arc<AppContext>>,
|
||||||
|
#[cfg(test)]
|
||||||
|
concurrency_manager: Option<Arc<ConcurrencyManager>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DefaultMultipartUsecase {
|
impl DefaultMultipartUsecase {
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pub fn without_context() -> Self {
|
pub fn without_context() -> Self {
|
||||||
Self { context: None }
|
Self {
|
||||||
|
context: None,
|
||||||
|
concurrency_manager: None,
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn from_global() -> Self {
|
pub fn from_global() -> Self {
|
||||||
Self {
|
Self {
|
||||||
context: current_app_context(),
|
context: current_app_context(),
|
||||||
|
#[cfg(test)]
|
||||||
|
concurrency_manager: None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -401,7 +412,22 @@ impl DefaultMultipartUsecase {
|
|||||||
/// so the use-case resolves that server's store; `None` falls back to the
|
/// so the use-case resolves that server's store; `None` falls back to the
|
||||||
/// ambient default.
|
/// ambient default.
|
||||||
pub fn with_context(context: Option<std::sync::Arc<crate::runtime_sources::AppContext>>) -> Self {
|
pub fn with_context(context: Option<std::sync::Arc<crate::runtime_sources::AppContext>>) -> Self {
|
||||||
Self { context }
|
Self {
|
||||||
|
context,
|
||||||
|
#[cfg(test)]
|
||||||
|
concurrency_manager: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
fn with_context_and_concurrency_manager(
|
||||||
|
context: Option<std::sync::Arc<crate::runtime_sources::AppContext>>,
|
||||||
|
concurrency_manager: Arc<ConcurrencyManager>,
|
||||||
|
) -> Self {
|
||||||
|
Self {
|
||||||
|
context,
|
||||||
|
concurrency_manager: Some(concurrency_manager),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn bucket_metadata_sys(&self) -> Option<Arc<RwLock<metadata_sys::BucketMetadataSys>>> {
|
fn bucket_metadata_sys(&self) -> Option<Arc<RwLock<metadata_sys::BucketMetadataSys>>> {
|
||||||
@@ -416,6 +442,15 @@ impl DefaultMultipartUsecase {
|
|||||||
current_object_data_cache_for_context(self.context.as_deref())
|
current_object_data_cache_for_context(self.context.as_deref())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn concurrency_manager(&self) -> &ConcurrencyManager {
|
||||||
|
#[cfg(test)]
|
||||||
|
if let Some(concurrency_manager) = self.concurrency_manager.as_deref() {
|
||||||
|
return concurrency_manager;
|
||||||
|
}
|
||||||
|
|
||||||
|
get_concurrency_manager()
|
||||||
|
}
|
||||||
|
|
||||||
#[instrument(level = "debug", skip(self))]
|
#[instrument(level = "debug", skip(self))]
|
||||||
pub async fn execute_abort_multipart_upload(
|
pub async fn execute_abort_multipart_upload(
|
||||||
&self,
|
&self,
|
||||||
@@ -1071,6 +1106,25 @@ impl DefaultMultipartUsecase {
|
|||||||
{
|
{
|
||||||
return Err(S3Error::new(S3ErrorCode::EntityTooLarge));
|
return Err(S3Error::new(S3ErrorCode::EntityTooLarge));
|
||||||
}
|
}
|
||||||
|
let upload_part_admission = match self
|
||||||
|
.concurrency_manager()
|
||||||
|
.admit_multipart_part(size.unwrap_or(-1))
|
||||||
|
.await
|
||||||
|
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "foreground write admission closed"))?
|
||||||
|
{
|
||||||
|
ForegroundWriteAdmission::Disabled => None,
|
||||||
|
ForegroundWriteAdmission::Admitted(permit) => {
|
||||||
|
counter!("rustfs.upload_part.foreground_admission.total", "result" => "admitted").increment(1);
|
||||||
|
Some(permit)
|
||||||
|
}
|
||||||
|
ForegroundWriteAdmission::Rejected => {
|
||||||
|
counter!("rustfs.upload_part.foreground_admission.total", "result" => "rejected").increment(1);
|
||||||
|
return Err(S3Error::with_message(
|
||||||
|
S3ErrorCode::SlowDown,
|
||||||
|
"foreground write concurrency limit reached, please reduce your request rate",
|
||||||
|
));
|
||||||
|
}
|
||||||
|
};
|
||||||
if max_total_object_size.is_some() {
|
if max_total_object_size.is_some() {
|
||||||
let request_id = req
|
let request_id = req
|
||||||
.extensions
|
.extensions
|
||||||
@@ -1252,10 +1306,12 @@ impl DefaultMultipartUsecase {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let _upload_part_admission = upload_part_admission;
|
||||||
let info = store
|
let info = store
|
||||||
.put_object_part(&bucket, &key, &upload_id, part_id, &mut reader, &opts)
|
.put_object_part(&bucket, &key, &upload_id, part_id, &mut reader, &opts)
|
||||||
.await
|
.await
|
||||||
.map_err(ApiError::from)?;
|
.map_err(ApiError::from)?;
|
||||||
|
drop(_upload_part_admission);
|
||||||
|
|
||||||
let mut checksum_crc32 = input.checksum_crc32;
|
let mut checksum_crc32 = input.checksum_crc32;
|
||||||
let mut checksum_crc32c = input.checksum_crc32c;
|
let mut checksum_crc32c = input.checksum_crc32c;
|
||||||
@@ -2930,4 +2986,52 @@ mod tests {
|
|||||||
assert_eq!(err.message(), Some("partNumber must be between 1 and 10000"));
|
assert_eq!(err.message(), Some("partNumber must be between 1 and 10000"));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn execute_upload_part_rejects_when_foreground_write_admission_is_full() {
|
||||||
|
use crate::app::storage_api::test::contract::bucket::{BucketOperations, MakeBucketOptions};
|
||||||
|
|
||||||
|
let store = crate::app::gating_test_env::shared_gating_ecstore().await;
|
||||||
|
let ambient = crate::app::gating_test_env::shared_gating_ambient().await;
|
||||||
|
let context = Arc::new(AppContext::new(Arc::clone(&store), ambient.iam(), ambient.kms()));
|
||||||
|
let bucket = format!("upload-part-admission-{}", Uuid::new_v4().simple());
|
||||||
|
store
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create upload part admission test bucket");
|
||||||
|
let upload = store
|
||||||
|
.new_multipart_upload(&bucket, "object", &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create multipart upload");
|
||||||
|
let concurrency_manager = Arc::new(ConcurrencyManager::with_large_put_admission_for_test(
|
||||||
|
true,
|
||||||
|
1,
|
||||||
|
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
||||||
|
Duration::ZERO,
|
||||||
|
));
|
||||||
|
let held = concurrency_manager
|
||||||
|
.admit_multipart_part(1024)
|
||||||
|
.await
|
||||||
|
.expect("first multipart part should acquire the only permit");
|
||||||
|
let usecase =
|
||||||
|
DefaultMultipartUsecase::with_context_and_concurrency_manager(Some(context), Arc::clone(&concurrency_manager));
|
||||||
|
let body = StreamingBlob::wrap(futures::stream::pending::<Result<Bytes, std::io::Error>>());
|
||||||
|
let input = UploadPartInput::builder()
|
||||||
|
.bucket(bucket)
|
||||||
|
.key("object".to_string())
|
||||||
|
.upload_id(upload.upload_id)
|
||||||
|
.part_number(1)
|
||||||
|
.content_length(Some(1024))
|
||||||
|
.body(Some(body))
|
||||||
|
.build()
|
||||||
|
.expect("upload part input should build");
|
||||||
|
|
||||||
|
let err = usecase
|
||||||
|
.execute_upload_part(build_request(input, Method::PUT))
|
||||||
|
.await
|
||||||
|
.expect_err("full foreground write admission should reject before body ingest");
|
||||||
|
assert_eq!(err.code(), &S3ErrorCode::SlowDown);
|
||||||
|
drop(held);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -57,8 +57,8 @@ use super::storage_api::object_usecase::bucket::{
|
|||||||
versioning_sys::BucketVersioningSys,
|
versioning_sys::BucketVersioningSys,
|
||||||
};
|
};
|
||||||
use super::storage_api::object_usecase::compression::{MIN_DISK_COMPRESSIBLE_SIZE, is_disk_compressible};
|
use super::storage_api::object_usecase::compression::{MIN_DISK_COMPRESSIBLE_SIZE, is_disk_compressible};
|
||||||
use super::storage_api::object_usecase::concurrency::{
|
pub(crate) use super::storage_api::object_usecase::concurrency::{
|
||||||
self, ConcurrencyManager, DiskReadAdmission, GetObjectGuard, PutObjectAdmission, PutObjectGuard,
|
self, ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, PutObjectGuard,
|
||||||
get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
|
get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
|
||||||
};
|
};
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
|
|||||||
@@ -966,12 +966,12 @@ impl DefaultObjectUsecase {
|
|||||||
.await
|
.await
|
||||||
.map_err(|_| s3_error!(InternalError, "foreground write admission closed"))?
|
.map_err(|_| s3_error!(InternalError, "foreground write admission closed"))?
|
||||||
{
|
{
|
||||||
PutObjectAdmission::Disabled => None,
|
ForegroundWriteAdmission::Disabled => None,
|
||||||
PutObjectAdmission::Admitted(permit) => {
|
ForegroundWriteAdmission::Admitted(permit) => {
|
||||||
counter!("rustfs.put_object.foreground_admission.total", "result" => "admitted").increment(1);
|
counter!("rustfs.put_object.foreground_admission.total", "result" => "admitted").increment(1);
|
||||||
Some(permit)
|
Some(permit)
|
||||||
}
|
}
|
||||||
PutObjectAdmission::Rejected => {
|
ForegroundWriteAdmission::Rejected => {
|
||||||
counter!("rustfs.put_object.foreground_admission.total", "result" => "rejected").increment(1);
|
counter!("rustfs.put_object.foreground_admission.total", "result" => "rejected").increment(1);
|
||||||
return Err(s3_error!(
|
return Err(s3_error!(
|
||||||
SlowDown,
|
SlowDown,
|
||||||
|
|||||||
@@ -979,8 +979,8 @@ pub(crate) mod bucket {
|
|||||||
|
|
||||||
pub(crate) mod concurrency {
|
pub(crate) mod concurrency {
|
||||||
pub(crate) use crate::storage::storage_api::concurrency_consumer::{
|
pub(crate) use crate::storage::storage_api::concurrency_consumer::{
|
||||||
ConcurrencyManager, DiskReadAdmission, GetObjectGuard, IoQueueStatus, IoStrategy, PutObjectAdmission, PutObjectGuard,
|
ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, IoQueueStatus, IoStrategy,
|
||||||
get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
|
PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -67,8 +67,8 @@ pub struct ConcurrencyManager {
|
|||||||
bandwidth_monitor: Arc<Mutex<BandwidthMonitor>>,
|
bandwidth_monitor: Arc<Mutex<BandwidthMonitor>>,
|
||||||
/// Metrics collector for I/O latency tracking (P50, P95, P99)
|
/// Metrics collector for I/O latency tracking (P50, P95, P99)
|
||||||
metrics_collector: Arc<MetricsCollector>,
|
metrics_collector: Arc<MetricsCollector>,
|
||||||
/// Foreground PutObject admission policy, resolved once at startup.
|
/// Foreground write admission policy, resolved once at startup.
|
||||||
put_admission_policy: PutAdmissionPolicy,
|
foreground_write_admission_policy: ForegroundWriteAdmissionPolicy,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl std::fmt::Debug for ConcurrencyManager {
|
impl std::fmt::Debug for ConcurrencyManager {
|
||||||
@@ -118,26 +118,26 @@ pub enum DiskReadAdmission {
|
|||||||
Rejected,
|
Rejected,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Outcome of foreground PutObject request admission.
|
/// Outcome of foreground write request admission.
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub enum PutObjectAdmission {
|
pub enum ForegroundWriteAdmission {
|
||||||
/// Foreground PUT admission is disabled; proceed on the legacy path.
|
/// Foreground write admission is disabled; proceed on the legacy path.
|
||||||
Disabled,
|
Disabled,
|
||||||
/// Request is admitted and must hold the permit until the store write
|
/// Request is admitted and must hold the permit until the store write
|
||||||
/// returns or the request fails before mutation.
|
/// returns or the request fails before mutation.
|
||||||
Admitted(tokio::sync::OwnedSemaphorePermit),
|
Admitted(tokio::sync::OwnedSemaphorePermit),
|
||||||
/// The selected foreground PUT admission gate stayed full until the configured wait timeout.
|
/// The selected foreground write admission gate stayed full until the configured wait timeout.
|
||||||
Rejected,
|
Rejected,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
struct PutAdmissionGate {
|
struct ForegroundWriteAdmissionGate {
|
||||||
semaphore: Arc<Semaphore>,
|
semaphore: Arc<Semaphore>,
|
||||||
limit: usize,
|
limit: usize,
|
||||||
wait_timeout: Duration,
|
wait_timeout: Duration,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl PutAdmissionGate {
|
impl ForegroundWriteAdmissionGate {
|
||||||
fn new(limit: usize, wait_timeout: Duration) -> Self {
|
fn new(limit: usize, wait_timeout: Duration) -> Self {
|
||||||
Self {
|
Self {
|
||||||
semaphore: Arc::new(Semaphore::new(limit)),
|
semaphore: Arc::new(Semaphore::new(limit)),
|
||||||
@@ -150,36 +150,46 @@ impl PutAdmissionGate {
|
|||||||
self.limit.saturating_sub(self.semaphore.available_permits())
|
self.limit.saturating_sub(self.semaphore.available_permits())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn admit(&self) -> Result<PutObjectAdmission, tokio::sync::AcquireError> {
|
async fn admit(&self) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
|
||||||
if self.wait_timeout.is_zero() {
|
if self.wait_timeout.is_zero() {
|
||||||
return Ok(match self.semaphore.clone().try_acquire_owned() {
|
return Ok(match self.semaphore.clone().try_acquire_owned() {
|
||||||
Ok(permit) => PutObjectAdmission::Admitted(permit),
|
Ok(permit) => ForegroundWriteAdmission::Admitted(permit),
|
||||||
Err(tokio::sync::TryAcquireError::NoPermits) => PutObjectAdmission::Rejected,
|
Err(tokio::sync::TryAcquireError::NoPermits) => ForegroundWriteAdmission::Rejected,
|
||||||
Err(tokio::sync::TryAcquireError::Closed) => PutObjectAdmission::Rejected,
|
Err(tokio::sync::TryAcquireError::Closed) => ForegroundWriteAdmission::Rejected,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
match tokio::time::timeout(self.wait_timeout, self.semaphore.clone().acquire_owned()).await {
|
match tokio::time::timeout(self.wait_timeout, self.semaphore.clone().acquire_owned()).await {
|
||||||
Ok(permit) => Ok(PutObjectAdmission::Admitted(permit?)),
|
Ok(permit) => Ok(ForegroundWriteAdmission::Admitted(permit?)),
|
||||||
Err(_) => Ok(PutObjectAdmission::Rejected),
|
Err(_) => Ok(ForegroundWriteAdmission::Rejected),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
enum PutAdmissionPolicy {
|
enum ForegroundWriteAdmissionPolicy {
|
||||||
/// Strict admission was explicitly enabled with limit `0`.
|
/// Strict admission was explicitly enabled with limit `0`.
|
||||||
Disabled,
|
Disabled,
|
||||||
/// No hard PUT gate is configured; foreground write snapshots use the
|
/// No hard PUT gate is configured; foreground write snapshots use the
|
||||||
/// existing active request counter as a soft pressure signal.
|
/// existing active request counter as a soft pressure signal.
|
||||||
LegacyCounterOnly,
|
LegacyCounterOnly,
|
||||||
/// Explicit all-PUT admission gate.
|
/// Explicit all foreground write admission gate.
|
||||||
Strict(PutAdmissionGate),
|
Strict(ForegroundWriteAdmissionGate),
|
||||||
/// Default large/unknown-size PUT admission gate.
|
/// Default foreground write admission gate for pressure-heavy writes.
|
||||||
Large { gate: PutAdmissionGate, min_size_bytes: usize },
|
Large {
|
||||||
|
gate: ForegroundWriteAdmissionGate,
|
||||||
|
put_object_min_size_bytes: usize,
|
||||||
|
multipart_part_min_size_bytes: usize,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
impl PutAdmissionPolicy {
|
#[derive(Clone, Copy)]
|
||||||
|
enum ForegroundWriteAdmissionKind {
|
||||||
|
PutObject,
|
||||||
|
MultipartPart,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ForegroundWriteAdmissionPolicy {
|
||||||
fn from_env(max_disk_reads: usize) -> Self {
|
fn from_env(max_disk_reads: usize) -> Self {
|
||||||
let strict_enabled = rustfs_utils::get_env_bool(
|
let strict_enabled = rustfs_utils::get_env_bool(
|
||||||
rustfs_config::ENV_PUT_FOREGROUND_ADMISSION_ENABLE,
|
rustfs_config::ENV_PUT_FOREGROUND_ADMISSION_ENABLE,
|
||||||
@@ -197,7 +207,7 @@ impl PutAdmissionPolicy {
|
|||||||
return if strict_limit == 0 {
|
return if strict_limit == 0 {
|
||||||
Self::Disabled
|
Self::Disabled
|
||||||
} else {
|
} else {
|
||||||
Self::Strict(PutAdmissionGate::new(strict_limit, strict_wait_timeout))
|
Self::Strict(ForegroundWriteAdmissionGate::new(strict_limit, strict_wait_timeout))
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -216,18 +226,23 @@ impl PutAdmissionPolicy {
|
|||||||
),
|
),
|
||||||
max_disk_reads,
|
max_disk_reads,
|
||||||
);
|
);
|
||||||
let min_size_bytes = rustfs_utils::get_env_usize(
|
let put_object_min_size_bytes = rustfs_utils::get_env_usize(
|
||||||
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
||||||
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
||||||
);
|
);
|
||||||
|
let multipart_part_min_size_bytes = rustfs_utils::get_env_usize(
|
||||||
|
rustfs_config::ENV_PUT_MULTIPART_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
||||||
|
rustfs_config::DEFAULT_PUT_MULTIPART_FOREGROUND_ADMISSION_MIN_SIZE_BYTES,
|
||||||
|
);
|
||||||
let wait_timeout = Duration::from_millis(rustfs_utils::get_env_u64(
|
let wait_timeout = Duration::from_millis(rustfs_utils::get_env_u64(
|
||||||
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS,
|
rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS,
|
||||||
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS,
|
rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS,
|
||||||
));
|
));
|
||||||
|
|
||||||
Self::Large {
|
Self::Large {
|
||||||
gate: PutAdmissionGate::new(large_limit, wait_timeout),
|
gate: ForegroundWriteAdmissionGate::new(large_limit, wait_timeout),
|
||||||
min_size_bytes,
|
put_object_min_size_bytes,
|
||||||
|
multipart_part_min_size_bytes,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -237,7 +252,7 @@ impl PutAdmissionPolicy {
|
|||||||
if limit == 0 {
|
if limit == 0 {
|
||||||
Self::Disabled
|
Self::Disabled
|
||||||
} else {
|
} else {
|
||||||
Self::Strict(PutAdmissionGate::new(limit, wait_timeout))
|
Self::Strict(ForegroundWriteAdmissionGate::new(limit, wait_timeout))
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
Self::LegacyCounterOnly
|
Self::LegacyCounterOnly
|
||||||
@@ -248,20 +263,38 @@ impl PutAdmissionPolicy {
|
|||||||
fn large_for_test(enabled: bool, limit: usize, min_size_bytes: usize, wait_timeout: Duration) -> Self {
|
fn large_for_test(enabled: bool, limit: usize, min_size_bytes: usize, wait_timeout: Duration) -> Self {
|
||||||
if enabled && limit > 0 {
|
if enabled && limit > 0 {
|
||||||
Self::Large {
|
Self::Large {
|
||||||
gate: PutAdmissionGate::new(limit, wait_timeout),
|
gate: ForegroundWriteAdmissionGate::new(limit, wait_timeout),
|
||||||
min_size_bytes,
|
put_object_min_size_bytes: min_size_bytes,
|
||||||
|
multipart_part_min_size_bytes: 0,
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
Self::LegacyCounterOnly
|
Self::LegacyCounterOnly
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn admit(&self, size: i64) -> Result<PutObjectAdmission, tokio::sync::AcquireError> {
|
async fn admit(
|
||||||
|
&self,
|
||||||
|
kind: ForegroundWriteAdmissionKind,
|
||||||
|
size: i64,
|
||||||
|
) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
|
||||||
match self {
|
match self {
|
||||||
Self::Disabled | Self::LegacyCounterOnly => Ok(PutObjectAdmission::Disabled),
|
Self::Disabled | Self::LegacyCounterOnly => Ok(ForegroundWriteAdmission::Disabled),
|
||||||
Self::Strict(gate) => gate.admit().await,
|
Self::Strict(gate) => gate.admit().await,
|
||||||
Self::Large { gate, min_size_bytes } if should_gate_large_put(size, *min_size_bytes) => gate.admit().await,
|
Self::Large {
|
||||||
Self::Large { .. } => Ok(PutObjectAdmission::Disabled),
|
gate,
|
||||||
|
put_object_min_size_bytes,
|
||||||
|
multipart_part_min_size_bytes,
|
||||||
|
} => {
|
||||||
|
let min_size_bytes = match kind {
|
||||||
|
ForegroundWriteAdmissionKind::PutObject => *put_object_min_size_bytes,
|
||||||
|
ForegroundWriteAdmissionKind::MultipartPart => *multipart_part_min_size_bytes,
|
||||||
|
};
|
||||||
|
if should_gate_foreground_write(size, min_size_bytes) {
|
||||||
|
gate.admit().await
|
||||||
|
} else {
|
||||||
|
Ok(ForegroundWriteAdmission::Disabled)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -313,7 +346,7 @@ fn derive_large_put_admission_limit(configured_limit: usize, max_disk_reads: usi
|
|||||||
scheduler_base.div_ceil(2).clamp(1, DERIVED_LARGE_PUT_ADMISSION_LIMIT_MAX)
|
scheduler_base.div_ceil(2).clamp(1, DERIVED_LARGE_PUT_ADMISSION_LIMIT_MAX)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn should_gate_large_put(size: i64, min_size_bytes: usize) -> bool {
|
fn should_gate_foreground_write(size: i64, min_size_bytes: usize) -> bool {
|
||||||
if min_size_bytes == 0 || size < 0 {
|
if min_size_bytes == 0 || size < 0 {
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
@@ -367,7 +400,7 @@ impl ConcurrencyManager {
|
|||||||
// Initialize metrics collector for I/O latency tracking
|
// Initialize metrics collector for I/O latency tracking
|
||||||
// Keep 1000 samples for P95/P99 calculation
|
// Keep 1000 samples for P95/P99 calculation
|
||||||
let metrics_collector = Arc::new(MetricsCollector::new(performance_metrics, 1000));
|
let metrics_collector = Arc::new(MetricsCollector::new(performance_metrics, 1000));
|
||||||
let put_admission_policy = PutAdmissionPolicy::from_env(max_disk_reads);
|
let foreground_write_admission_policy = ForegroundWriteAdmissionPolicy::from_env(max_disk_reads);
|
||||||
|
|
||||||
// Build queue config directly from scheduler config.
|
// Build queue config directly from scheduler config.
|
||||||
let queue_config = IoPriorityQueueConfig::from_scheduler_config(&scheduler_config);
|
let queue_config = IoPriorityQueueConfig::from_scheduler_config(&scheduler_config);
|
||||||
@@ -383,7 +416,7 @@ impl ConcurrencyManager {
|
|||||||
pattern_detector,
|
pattern_detector,
|
||||||
bandwidth_monitor,
|
bandwidth_monitor,
|
||||||
metrics_collector,
|
metrics_collector,
|
||||||
put_admission_policy,
|
foreground_write_admission_policy,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -410,7 +443,7 @@ impl ConcurrencyManager {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pub(crate) fn with_put_admission_for_test(enabled: bool, limit: usize, wait_timeout: Duration) -> Self {
|
pub(crate) fn with_put_admission_for_test(enabled: bool, limit: usize, wait_timeout: Duration) -> Self {
|
||||||
let mut manager = Self::new();
|
let mut manager = Self::new();
|
||||||
manager.put_admission_policy = PutAdmissionPolicy::strict_for_test(enabled, limit, wait_timeout);
|
manager.foreground_write_admission_policy = ForegroundWriteAdmissionPolicy::strict_for_test(enabled, limit, wait_timeout);
|
||||||
manager
|
manager
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -422,7 +455,8 @@ impl ConcurrencyManager {
|
|||||||
wait_timeout: Duration,
|
wait_timeout: Duration,
|
||||||
) -> Self {
|
) -> Self {
|
||||||
let mut manager = Self::new();
|
let mut manager = Self::new();
|
||||||
manager.put_admission_policy = PutAdmissionPolicy::large_for_test(enabled, limit, min_size_bytes, wait_timeout);
|
manager.foreground_write_admission_policy =
|
||||||
|
ForegroundWriteAdmissionPolicy::large_for_test(enabled, limit, min_size_bytes, wait_timeout);
|
||||||
manager
|
manager
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -516,8 +550,21 @@ impl ConcurrencyManager {
|
|||||||
/// The strict experimental gate applies to every PUT only when explicitly
|
/// The strict experimental gate applies to every PUT only when explicitly
|
||||||
/// enabled. Otherwise the default-on large-object gate protects sustained
|
/// enabled. Otherwise the default-on large-object gate protects sustained
|
||||||
/// erasure/RPC pressure while keeping small PUTs on the legacy path.
|
/// erasure/RPC pressure while keeping small PUTs on the legacy path.
|
||||||
pub async fn admit_put_object(&self, size: i64) -> Result<PutObjectAdmission, tokio::sync::AcquireError> {
|
pub async fn admit_put_object(&self, size: i64) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
|
||||||
self.put_admission_policy.admit(size).await
|
self.foreground_write_admission_policy
|
||||||
|
.admit(ForegroundWriteAdmissionKind::PutObject, size)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Admit a multipart UploadPart request under the configured write gate.
|
||||||
|
///
|
||||||
|
/// Multipart workloads can saturate memory and internode write streams with
|
||||||
|
/// many moderate-sized parts, so they use an operation-specific threshold
|
||||||
|
/// while sharing the same foreground write permit pool.
|
||||||
|
pub async fn admit_multipart_part(&self, size: i64) -> Result<ForegroundWriteAdmission, tokio::sync::AcquireError> {
|
||||||
|
self.foreground_write_admission_policy
|
||||||
|
.admit(ForegroundWriteAdmissionKind::MultipartPart, size)
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
// ============================================
|
// ============================================
|
||||||
@@ -928,7 +975,8 @@ impl ConcurrencyManager {
|
|||||||
|
|
||||||
/// Get a read-only workload admission snapshot for foreground writes.
|
/// Get a read-only workload admission snapshot for foreground writes.
|
||||||
pub fn put_object_admission_snapshot(&self) -> WorkloadAdmissionSnapshot {
|
pub fn put_object_admission_snapshot(&self) -> WorkloadAdmissionSnapshot {
|
||||||
self.put_admission_policy.snapshot(self.scheduler_config.max_concurrent_reads)
|
self.foreground_write_admission_policy
|
||||||
|
.snapshot(self.scheduler_config.max_concurrent_reads)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Get a read-only workload admission registry snapshot for local storage concurrency.
|
/// Get a read-only workload admission registry snapshot for local storage concurrency.
|
||||||
@@ -1002,7 +1050,7 @@ impl Default for ConcurrencyManager {
|
|||||||
mod integration_tests {
|
mod integration_tests {
|
||||||
use super::super::io_schedule::{IoLoadLevel, IoPriority};
|
use super::super::io_schedule::{IoLoadLevel, IoPriority};
|
||||||
use super::super::request_guard::GetObjectGuard;
|
use super::super::request_guard::GetObjectGuard;
|
||||||
use super::{ConcurrencyManager, PutObjectAdmission, derive_large_put_admission_limit};
|
use super::{ConcurrencyManager, ForegroundWriteAdmission, derive_large_put_admission_limit};
|
||||||
use crate::storage::storage_api::concurrency_consumer::PutObjectGuard;
|
use crate::storage::storage_api::concurrency_consumer::PutObjectGuard;
|
||||||
use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass};
|
use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass};
|
||||||
use rustfs_io_core::io_profile::{AccessPattern, StorageMedia};
|
use rustfs_io_core::io_profile::{AccessPattern, StorageMedia};
|
||||||
@@ -1109,7 +1157,7 @@ mod integration_tests {
|
|||||||
.await
|
.await
|
||||||
.expect("disabled put admission must not close");
|
.expect("disabled put admission must not close");
|
||||||
|
|
||||||
assert!(matches!(admission, PutObjectAdmission::Disabled));
|
assert!(matches!(admission, ForegroundWriteAdmission::Disabled));
|
||||||
assert_eq!(manager.put_object_admission_snapshot().state, AdmissionState::Open);
|
assert_eq!(manager.put_object_admission_snapshot().state, AdmissionState::Open);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1123,7 +1171,7 @@ mod integration_tests {
|
|||||||
.await
|
.await
|
||||||
.expect("strict zero-limit put admission must not close");
|
.expect("strict zero-limit put admission must not close");
|
||||||
|
|
||||||
assert!(matches!(admission, PutObjectAdmission::Disabled));
|
assert!(matches!(admission, ForegroundWriteAdmission::Disabled));
|
||||||
assert_eq!(manager.put_object_admission_snapshot().state, AdmissionState::Disabled);
|
assert_eq!(manager.put_object_admission_snapshot().state, AdmissionState::Disabled);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1136,14 +1184,14 @@ mod integration_tests {
|
|||||||
.admit_put_object(1024)
|
.admit_put_object(1024)
|
||||||
.await
|
.await
|
||||||
.expect("first put admission should acquire");
|
.expect("first put admission should acquire");
|
||||||
assert!(matches!(first, PutObjectAdmission::Admitted(_)));
|
assert!(matches!(first, ForegroundWriteAdmission::Admitted(_)));
|
||||||
assert_eq!(manager.put_object_admission_snapshot().state, AdmissionState::Saturated);
|
assert_eq!(manager.put_object_admission_snapshot().state, AdmissionState::Saturated);
|
||||||
|
|
||||||
let second = manager
|
let second = manager
|
||||||
.admit_put_object(1024)
|
.admit_put_object(1024)
|
||||||
.await
|
.await
|
||||||
.expect("full put admission gate should reject, not close");
|
.expect("full put admission gate should reject, not close");
|
||||||
assert!(matches!(second, PutObjectAdmission::Rejected));
|
assert!(matches!(second, ForegroundWriteAdmission::Rejected));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -1161,7 +1209,7 @@ mod integration_tests {
|
|||||||
.admit_put_object(1024)
|
.admit_put_object(1024)
|
||||||
.await
|
.await
|
||||||
.expect("released put admission permit should be reusable");
|
.expect("released put admission permit should be reusable");
|
||||||
assert!(matches!(second, PutObjectAdmission::Admitted(_)));
|
assert!(matches!(second, ForegroundWriteAdmission::Admitted(_)));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test(start_paused = true)]
|
#[tokio::test(start_paused = true)]
|
||||||
@@ -1182,7 +1230,7 @@ mod integration_tests {
|
|||||||
.await
|
.await
|
||||||
.expect("put admission waiter task must not panic")
|
.expect("put admission waiter task must not panic")
|
||||||
.expect("put admission gate must stay open");
|
.expect("put admission gate must stay open");
|
||||||
assert!(matches!(admission, PutObjectAdmission::Rejected));
|
assert!(matches!(admission, ForegroundWriteAdmission::Rejected));
|
||||||
drop(held);
|
drop(held);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1196,19 +1244,19 @@ mod integration_tests {
|
|||||||
.admit_put_object(min_size as i64)
|
.admit_put_object(min_size as i64)
|
||||||
.await
|
.await
|
||||||
.expect("large put admission should acquire");
|
.expect("large put admission should acquire");
|
||||||
assert!(matches!(held, PutObjectAdmission::Admitted(_)));
|
assert!(matches!(held, ForegroundWriteAdmission::Admitted(_)));
|
||||||
|
|
||||||
let small = manager
|
let small = manager
|
||||||
.admit_put_object((min_size - 1) as i64)
|
.admit_put_object((min_size - 1) as i64)
|
||||||
.await
|
.await
|
||||||
.expect("small put should bypass large admission");
|
.expect("small put should bypass large admission");
|
||||||
assert!(matches!(small, PutObjectAdmission::Disabled));
|
assert!(matches!(small, ForegroundWriteAdmission::Disabled));
|
||||||
|
|
||||||
let large = manager
|
let large = manager
|
||||||
.admit_put_object(min_size as i64)
|
.admit_put_object(min_size as i64)
|
||||||
.await
|
.await
|
||||||
.expect("second large put should reject when the gate is full");
|
.expect("second large put should reject when the gate is full");
|
||||||
assert!(matches!(large, PutObjectAdmission::Rejected));
|
assert!(matches!(large, ForegroundWriteAdmission::Rejected));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -1220,13 +1268,38 @@ mod integration_tests {
|
|||||||
.admit_put_object(-1)
|
.admit_put_object(-1)
|
||||||
.await
|
.await
|
||||||
.expect("unknown-size put admission should acquire");
|
.expect("unknown-size put admission should acquire");
|
||||||
assert!(matches!(held, PutObjectAdmission::Admitted(_)));
|
assert!(matches!(held, ForegroundWriteAdmission::Admitted(_)));
|
||||||
|
|
||||||
let second = manager
|
let second = manager
|
||||||
.admit_put_object(-1)
|
.admit_put_object(-1)
|
||||||
.await
|
.await
|
||||||
.expect("unknown-size put admission should reject when the gate is full");
|
.expect("unknown-size put admission should reject when the gate is full");
|
||||||
assert!(matches!(second, PutObjectAdmission::Rejected));
|
assert!(matches!(second, ForegroundWriteAdmission::Rejected));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn test_concurrency_manager_large_put_admission_gates_multipart_parts_by_default() {
|
||||||
|
let min_size = rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_MIN_SIZE_BYTES;
|
||||||
|
let manager = ConcurrencyManager::with_large_put_admission_for_test(true, 1, min_size, Duration::ZERO);
|
||||||
|
|
||||||
|
let held = manager
|
||||||
|
.admit_multipart_part(1024)
|
||||||
|
.await
|
||||||
|
.expect("first multipart part admission should acquire");
|
||||||
|
assert!(matches!(held, ForegroundWriteAdmission::Admitted(_)));
|
||||||
|
|
||||||
|
let direct_small_put = manager
|
||||||
|
.admit_put_object((min_size - 1) as i64)
|
||||||
|
.await
|
||||||
|
.expect("small direct put should bypass the large put threshold");
|
||||||
|
assert!(matches!(direct_small_put, ForegroundWriteAdmission::Disabled));
|
||||||
|
|
||||||
|
let second_part = manager
|
||||||
|
.admit_multipart_part(1024)
|
||||||
|
.await
|
||||||
|
.expect("full multipart admission gate should reject, not close");
|
||||||
|
assert!(matches!(second_part, ForegroundWriteAdmission::Rejected));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
|||||||
@@ -51,7 +51,7 @@ pub use io_schedule::{
|
|||||||
pub use request_guard::{GetObjectGuard, PutObjectGuard};
|
pub use request_guard::{GetObjectGuard, PutObjectGuard};
|
||||||
|
|
||||||
// Concurrency manager
|
// Concurrency manager
|
||||||
pub use manager::{ConcurrencyManager, DiskReadAdmission, PutObjectAdmission};
|
pub use manager::{ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission};
|
||||||
|
|
||||||
// ============================================
|
// ============================================
|
||||||
// Helper Functions
|
// Helper Functions
|
||||||
|
|||||||
@@ -126,8 +126,8 @@ pub(crate) mod access_consumer {
|
|||||||
|
|
||||||
pub(crate) mod concurrency_consumer {
|
pub(crate) mod concurrency_consumer {
|
||||||
pub(crate) use super::super::concurrency::{
|
pub(crate) use super::super::concurrency::{
|
||||||
ConcurrencyManager, DiskReadAdmission, GetObjectGuard, IoQueueStatus, IoStrategy, PutObjectAdmission, PutObjectGuard,
|
ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, IoQueueStatus, IoStrategy,
|
||||||
get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
|
PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -44,7 +44,11 @@ cd "$(dirname "$0")/.."
|
|||||||
# inlined one s3_error! call in transport.rs (+1); measured 1615 on the
|
# inlined one s3_error! call in transport.rs (+1); measured 1615 on the
|
||||||
# pre-move main (after #6694) and 1616 after, so the slack 1620 baseline is
|
# pre-move main (after #6694) and 1616 after, so the slack 1620 baseline is
|
||||||
# retightened to the measured 1616.
|
# retightened to the measured 1616.
|
||||||
S3S_IMPORT_FILES_BASELINE=215
|
# 215 → 213 on 2026-08-28: multipart foreground admission cleanup moved the
|
||||||
|
# new usecase dependency behind the app object-domain facade while main had
|
||||||
|
# already shed two direct s3s-importing files. Retighten the file counter only;
|
||||||
|
# s3_error! stays flat at 1616.
|
||||||
|
S3S_IMPORT_FILES_BASELINE=213
|
||||||
S3_ERROR_LINES_BASELINE=1616
|
S3_ERROR_LINES_BASELINE=1616
|
||||||
# ecstore-scoped ratchet (rustfs/backlog#1842): the storage engine must not
|
# ecstore-scoped ratchet (rustfs/backlog#1842): the storage engine must not
|
||||||
# know S3 wire/DTO types (ARCHITECTURE.md invariant 4). The S3-*consuming*
|
# know S3 wire/DTO types (ARCHITECTURE.md invariant 4). The S3-*consuming*
|
||||||
|
|||||||
Reference in New Issue
Block a user