diff --git a/crates/config/src/constants/object.rs b/crates/config/src/constants/object.rs index 095ea2b3b..f1ead7972 100644 --- a/crates/config/src/constants/object.rs +++ b/crates/config/src/constants/object.rs @@ -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 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, /// 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 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 /// 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 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 /// `SlowDown`/503 before body ingest when the node is already saturated. diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index 557f2d834..86acb2039 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -67,7 +67,10 @@ use super::storage_api::multipart_usecase::sse::{ use super::storage_api::multipart_usecase::{ 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::{ ObjectDataCacheAdapter, invalidate_object_data_cache_after_complete_multipart_success, invalidate_object_data_cache_before_mutation, @@ -90,6 +93,7 @@ use crate::table_catalog; use bytes::Bytes; use futures::StreamExt; use http::{HeaderMap, HeaderValue, Uri}; +use metrics::counter; use rustfs_io_metrics::record_s3_op; use rustfs_s3_ops::S3Operation; use rustfs_targets::EventName; @@ -382,17 +386,24 @@ fn build_complete_multipart_location(headers: &HeaderMap, uri: &Uri, bucket: &st #[derive(Clone, Default)] pub struct DefaultMultipartUsecase { context: Option>, + #[cfg(test)] + concurrency_manager: Option>, } impl DefaultMultipartUsecase { #[cfg(test)] pub fn without_context() -> Self { - Self { context: None } + Self { + context: None, + concurrency_manager: None, + } } pub fn from_global() -> Self { Self { 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 /// ambient default. pub fn with_context(context: Option>) -> Self { - Self { context } + Self { + context, + #[cfg(test)] + concurrency_manager: None, + } + } + + #[cfg(test)] + fn with_context_and_concurrency_manager( + context: Option>, + concurrency_manager: Arc, + ) -> Self { + Self { + context, + concurrency_manager: Some(concurrency_manager), + } } fn bucket_metadata_sys(&self) -> Option>> { @@ -416,6 +442,15 @@ impl DefaultMultipartUsecase { 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))] pub async fn execute_abort_multipart_upload( &self, @@ -1071,6 +1106,25 @@ impl DefaultMultipartUsecase { { 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() { let request_id = req .extensions @@ -1252,10 +1306,12 @@ impl DefaultMultipartUsecase { ); } + let _upload_part_admission = upload_part_admission; let info = store .put_object_part(&bucket, &key, &upload_id, part_id, &mut reader, &opts) .await .map_err(ApiError::from)?; + drop(_upload_part_admission); let mut checksum_crc32 = input.checksum_crc32; 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")); } } + + #[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::>()); + 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); + } } diff --git a/rustfs/src/app/object/mod.rs b/rustfs/src/app/object/mod.rs index 1e42d7036..eaf083ee7 100644 --- a/rustfs/src/app/object/mod.rs +++ b/rustfs/src/app/object/mod.rs @@ -57,8 +57,8 @@ use super::storage_api::object_usecase::bucket::{ versioning_sys::BucketVersioningSys, }; use super::storage_api::object_usecase::compression::{MIN_DISK_COMPRESSIBLE_SIZE, is_disk_compressible}; -use super::storage_api::object_usecase::concurrency::{ - self, ConcurrencyManager, DiskReadAdmission, GetObjectGuard, PutObjectAdmission, PutObjectGuard, +pub(crate) use super::storage_api::object_usecase::concurrency::{ + self, ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size, }; #[cfg(test)] diff --git a/rustfs/src/app/object/put.rs b/rustfs/src/app/object/put.rs index a7211d8ef..386089cc1 100644 --- a/rustfs/src/app/object/put.rs +++ b/rustfs/src/app/object/put.rs @@ -966,12 +966,12 @@ impl DefaultObjectUsecase { .await .map_err(|_| s3_error!(InternalError, "foreground write admission closed"))? { - PutObjectAdmission::Disabled => None, - PutObjectAdmission::Admitted(permit) => { + ForegroundWriteAdmission::Disabled => None, + ForegroundWriteAdmission::Admitted(permit) => { counter!("rustfs.put_object.foreground_admission.total", "result" => "admitted").increment(1); Some(permit) } - PutObjectAdmission::Rejected => { + ForegroundWriteAdmission::Rejected => { counter!("rustfs.put_object.foreground_admission.total", "result" => "rejected").increment(1); return Err(s3_error!( SlowDown, diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index 85967ab55..da3c2fb1f 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -979,8 +979,8 @@ pub(crate) mod bucket { pub(crate) mod concurrency { pub(crate) use crate::storage::storage_api::concurrency_consumer::{ - ConcurrencyManager, DiskReadAdmission, GetObjectGuard, IoQueueStatus, IoStrategy, PutObjectAdmission, PutObjectGuard, - get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size, + ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, IoQueueStatus, IoStrategy, + PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size, }; } diff --git a/rustfs/src/storage/concurrency/manager.rs b/rustfs/src/storage/concurrency/manager.rs index 646991ff1..4c5d20ef0 100644 --- a/rustfs/src/storage/concurrency/manager.rs +++ b/rustfs/src/storage/concurrency/manager.rs @@ -67,8 +67,8 @@ pub struct ConcurrencyManager { bandwidth_monitor: Arc>, /// Metrics collector for I/O latency tracking (P50, P95, P99) metrics_collector: Arc, - /// Foreground PutObject admission policy, resolved once at startup. - put_admission_policy: PutAdmissionPolicy, + /// Foreground write admission policy, resolved once at startup. + foreground_write_admission_policy: ForegroundWriteAdmissionPolicy, } impl std::fmt::Debug for ConcurrencyManager { @@ -118,26 +118,26 @@ pub enum DiskReadAdmission { Rejected, } -/// Outcome of foreground PutObject request admission. +/// Outcome of foreground write request admission. #[derive(Debug)] -pub enum PutObjectAdmission { - /// Foreground PUT admission is disabled; proceed on the legacy path. +pub enum ForegroundWriteAdmission { + /// Foreground write admission is disabled; proceed on the legacy path. Disabled, /// Request is admitted and must hold the permit until the store write /// returns or the request fails before mutation. 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, } #[derive(Clone)] -struct PutAdmissionGate { +struct ForegroundWriteAdmissionGate { semaphore: Arc, limit: usize, wait_timeout: Duration, } -impl PutAdmissionGate { +impl ForegroundWriteAdmissionGate { fn new(limit: usize, wait_timeout: Duration) -> Self { Self { semaphore: Arc::new(Semaphore::new(limit)), @@ -150,36 +150,46 @@ impl PutAdmissionGate { self.limit.saturating_sub(self.semaphore.available_permits()) } - async fn admit(&self) -> Result { + async fn admit(&self) -> Result { if self.wait_timeout.is_zero() { return Ok(match self.semaphore.clone().try_acquire_owned() { - Ok(permit) => PutObjectAdmission::Admitted(permit), - Err(tokio::sync::TryAcquireError::NoPermits) => PutObjectAdmission::Rejected, - Err(tokio::sync::TryAcquireError::Closed) => PutObjectAdmission::Rejected, + Ok(permit) => ForegroundWriteAdmission::Admitted(permit), + Err(tokio::sync::TryAcquireError::NoPermits) => ForegroundWriteAdmission::Rejected, + Err(tokio::sync::TryAcquireError::Closed) => ForegroundWriteAdmission::Rejected, }); } match tokio::time::timeout(self.wait_timeout, self.semaphore.clone().acquire_owned()).await { - Ok(permit) => Ok(PutObjectAdmission::Admitted(permit?)), - Err(_) => Ok(PutObjectAdmission::Rejected), + Ok(permit) => Ok(ForegroundWriteAdmission::Admitted(permit?)), + Err(_) => Ok(ForegroundWriteAdmission::Rejected), } } } #[derive(Clone)] -enum PutAdmissionPolicy { +enum ForegroundWriteAdmissionPolicy { /// Strict admission was explicitly enabled with limit `0`. Disabled, /// No hard PUT gate is configured; foreground write snapshots use the /// existing active request counter as a soft pressure signal. LegacyCounterOnly, - /// Explicit all-PUT admission gate. - Strict(PutAdmissionGate), - /// Default large/unknown-size PUT admission gate. - Large { gate: PutAdmissionGate, min_size_bytes: usize }, + /// Explicit all foreground write admission gate. + Strict(ForegroundWriteAdmissionGate), + /// Default foreground write admission gate for pressure-heavy writes. + 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 { let strict_enabled = rustfs_utils::get_env_bool( rustfs_config::ENV_PUT_FOREGROUND_ADMISSION_ENABLE, @@ -197,7 +207,7 @@ impl PutAdmissionPolicy { return if strict_limit == 0 { Self::Disabled } 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, ); - 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::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( rustfs_config::ENV_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS, rustfs_config::DEFAULT_PUT_LARGE_FOREGROUND_ADMISSION_WAIT_TIMEOUT_MS, )); Self::Large { - gate: PutAdmissionGate::new(large_limit, wait_timeout), - min_size_bytes, + gate: ForegroundWriteAdmissionGate::new(large_limit, wait_timeout), + put_object_min_size_bytes, + multipart_part_min_size_bytes, } } @@ -237,7 +252,7 @@ impl PutAdmissionPolicy { if limit == 0 { Self::Disabled } else { - Self::Strict(PutAdmissionGate::new(limit, wait_timeout)) + Self::Strict(ForegroundWriteAdmissionGate::new(limit, wait_timeout)) } } else { Self::LegacyCounterOnly @@ -248,20 +263,38 @@ impl PutAdmissionPolicy { fn large_for_test(enabled: bool, limit: usize, min_size_bytes: usize, wait_timeout: Duration) -> Self { if enabled && limit > 0 { Self::Large { - gate: PutAdmissionGate::new(limit, wait_timeout), - min_size_bytes, + gate: ForegroundWriteAdmissionGate::new(limit, wait_timeout), + put_object_min_size_bytes: min_size_bytes, + multipart_part_min_size_bytes: 0, } } else { Self::LegacyCounterOnly } } - async fn admit(&self, size: i64) -> Result { + async fn admit( + &self, + kind: ForegroundWriteAdmissionKind, + size: i64, + ) -> Result { match self { - Self::Disabled | Self::LegacyCounterOnly => Ok(PutObjectAdmission::Disabled), + Self::Disabled | Self::LegacyCounterOnly => Ok(ForegroundWriteAdmission::Disabled), 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 { .. } => Ok(PutObjectAdmission::Disabled), + Self::Large { + 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) } -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 { return true; } @@ -367,7 +400,7 @@ impl ConcurrencyManager { // Initialize metrics collector for I/O latency tracking // Keep 1000 samples for P95/P99 calculation 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. let queue_config = IoPriorityQueueConfig::from_scheduler_config(&scheduler_config); @@ -383,7 +416,7 @@ impl ConcurrencyManager { pattern_detector, bandwidth_monitor, metrics_collector, - put_admission_policy, + foreground_write_admission_policy, } } @@ -410,7 +443,7 @@ impl ConcurrencyManager { #[cfg(test)] pub(crate) fn with_put_admission_for_test(enabled: bool, limit: usize, wait_timeout: Duration) -> Self { 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 } @@ -422,7 +455,8 @@ impl ConcurrencyManager { wait_timeout: Duration, ) -> Self { 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 } @@ -516,8 +550,21 @@ impl ConcurrencyManager { /// The strict experimental gate applies to every PUT only when explicitly /// enabled. Otherwise the default-on large-object gate protects sustained /// erasure/RPC pressure while keeping small PUTs on the legacy path. - pub async fn admit_put_object(&self, size: i64) -> Result { - self.put_admission_policy.admit(size).await + pub async fn admit_put_object(&self, size: i64) -> Result { + 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 { + 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. 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. @@ -1002,7 +1050,7 @@ impl Default for ConcurrencyManager { mod integration_tests { use super::super::io_schedule::{IoLoadLevel, IoPriority}; 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 rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass}; use rustfs_io_core::io_profile::{AccessPattern, StorageMedia}; @@ -1109,7 +1157,7 @@ mod integration_tests { .await .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); } @@ -1123,7 +1171,7 @@ mod integration_tests { .await .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); } @@ -1136,14 +1184,14 @@ mod integration_tests { .admit_put_object(1024) .await .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); let second = manager .admit_put_object(1024) .await .expect("full put admission gate should reject, not close"); - assert!(matches!(second, PutObjectAdmission::Rejected)); + assert!(matches!(second, ForegroundWriteAdmission::Rejected)); } #[tokio::test] @@ -1161,7 +1209,7 @@ mod integration_tests { .admit_put_object(1024) .await .expect("released put admission permit should be reusable"); - assert!(matches!(second, PutObjectAdmission::Admitted(_))); + assert!(matches!(second, ForegroundWriteAdmission::Admitted(_))); } #[tokio::test(start_paused = true)] @@ -1182,7 +1230,7 @@ mod integration_tests { .await .expect("put admission waiter task must not panic") .expect("put admission gate must stay open"); - assert!(matches!(admission, PutObjectAdmission::Rejected)); + assert!(matches!(admission, ForegroundWriteAdmission::Rejected)); drop(held); } @@ -1196,19 +1244,19 @@ mod integration_tests { .admit_put_object(min_size as i64) .await .expect("large put admission should acquire"); - assert!(matches!(held, PutObjectAdmission::Admitted(_))); + assert!(matches!(held, ForegroundWriteAdmission::Admitted(_))); let small = manager .admit_put_object((min_size - 1) as i64) .await .expect("small put should bypass large admission"); - assert!(matches!(small, PutObjectAdmission::Disabled)); + assert!(matches!(small, ForegroundWriteAdmission::Disabled)); let large = manager .admit_put_object(min_size as i64) .await .expect("second large put should reject when the gate is full"); - assert!(matches!(large, PutObjectAdmission::Rejected)); + assert!(matches!(large, ForegroundWriteAdmission::Rejected)); } #[tokio::test] @@ -1220,13 +1268,38 @@ mod integration_tests { .admit_put_object(-1) .await .expect("unknown-size put admission should acquire"); - assert!(matches!(held, PutObjectAdmission::Admitted(_))); + assert!(matches!(held, ForegroundWriteAdmission::Admitted(_))); let second = manager .admit_put_object(-1) .await .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] diff --git a/rustfs/src/storage/concurrency/mod.rs b/rustfs/src/storage/concurrency/mod.rs index eb47aa4ed..7996649ab 100644 --- a/rustfs/src/storage/concurrency/mod.rs +++ b/rustfs/src/storage/concurrency/mod.rs @@ -51,7 +51,7 @@ pub use io_schedule::{ pub use request_guard::{GetObjectGuard, PutObjectGuard}; // Concurrency manager -pub use manager::{ConcurrencyManager, DiskReadAdmission, PutObjectAdmission}; +pub use manager::{ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission}; // ============================================ // Helper Functions diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index acab79166..b70db13f5 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -126,8 +126,8 @@ pub(crate) mod access_consumer { pub(crate) mod concurrency_consumer { pub(crate) use super::super::concurrency::{ - ConcurrencyManager, DiskReadAdmission, GetObjectGuard, IoQueueStatus, IoStrategy, PutObjectAdmission, PutObjectGuard, - get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size, + ConcurrencyManager, DiskReadAdmission, ForegroundWriteAdmission, GetObjectGuard, IoQueueStatus, IoStrategy, + PutObjectGuard, get_concurrency_aware_buffer_size, get_concurrency_manager, get_put_concurrency_aware_buffer_size, }; } diff --git a/scripts/check_s3s_footprint.sh b/scripts/check_s3s_footprint.sh index 6132536ae..3dcb59a0b 100755 --- a/scripts/check_s3s_footprint.sh +++ b/scripts/check_s3s_footprint.sh @@ -44,7 +44,11 @@ cd "$(dirname "$0")/.." # 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 # 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 # ecstore-scoped ratchet (rustfs/backlog#1842): the storage engine must not # know S3 wire/DTO types (ARCHITECTURE.md invariant 4). The S3-*consuming*