From 30957e2b51ace86c98efb6d853b210820a0168c0 Mon Sep 17 00:00:00 2001 From: cxymds Date: Fri, 26 Jun 2026 17:48:09 +0800 Subject: [PATCH] fix(ecstore): throttle data movement under read pressure (#3906) --- Cargo.lock | 1 + crates/ecstore/Cargo.toml | 1 + .../ecstore/src/data_movement_backpressure.rs | 352 ++++++++++++++++++ crates/ecstore/src/lib.rs | 12 + crates/ecstore/src/pools.rs | 39 +- crates/ecstore/src/rebalance/entry.rs | 20 + crates/ecstore/src/runtime_sources.rs | 21 +- rustfs/src/startup_background.rs | 8 +- rustfs/src/storage/storage_api.rs | 7 + rustfs/src/storage_api.rs | 4 +- 10 files changed, 459 insertions(+), 6 deletions(-) create mode 100644 crates/ecstore/src/data_movement_backpressure.rs diff --git a/Cargo.lock b/Cargo.lock index 7c9c75fd4..c7ceb9d98 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9341,6 +9341,7 @@ dependencies = [ "rmp-serde", "rustfs-checksums", "rustfs-common", + "rustfs-concurrency", "rustfs-config", "rustfs-credentials", "rustfs-data-usage", diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index 9eaf91c8c..ee4b6694e 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -44,6 +44,7 @@ rustfs-storage-api.workspace = true rustfs-tls-runtime.workspace = true rustfs-checksums.workspace = true rustfs-config = { workspace = true, features = ["constants", "notify", "audit", "server-config-model"] } +rustfs-concurrency.workspace = true rustfs-credentials = { workspace = true } rustfs-common.workspace = true rustfs-policy.workspace = true diff --git a/crates/ecstore/src/data_movement_backpressure.rs b/crates/ecstore/src/data_movement_backpressure.rs new file mode 100644 index 000000000..84664c1b7 --- /dev/null +++ b/crates/ecstore/src/data_movement_backpressure.rs @@ -0,0 +1,352 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use crate::error::{Error, Result}; +use crate::runtime_sources::{self, WorkloadSnapshotProviderRef}; +use metrics::{counter, histogram}; +use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass}; +use std::time::{Duration, Instant}; +use tokio::time::sleep; +use tokio_util::sync::CancellationToken; +use tracing::debug; + +const LOG_COMPONENT_ECSTORE: &str = "ecstore"; +const LOG_SUBSYSTEM_DATA_MOVEMENT: &str = "data_movement"; +const EVENT_DATA_MOVEMENT_BACKPRESSURE: &str = "data_movement_backpressure"; +const DATA_MOVEMENT_BACKPRESSURE_ENABLE_ENV: &str = "RUSTFS_DATA_MOVEMENT_BACKPRESSURE_ENABLE"; +const DATA_MOVEMENT_FOREGROUND_READ_HIGH_PERCENT_ENV: &str = "RUSTFS_DATA_MOVEMENT_FOREGROUND_READ_HIGH_PERCENT"; +const DATA_MOVEMENT_RECHECK_MS_ENV: &str = "RUSTFS_DATA_MOVEMENT_RECHECK_MS"; +const DEFAULT_DATA_MOVEMENT_BACKPRESSURE_ENABLE: bool = true; +const DEFAULT_DATA_MOVEMENT_FOREGROUND_READ_HIGH_PERCENT: usize = 80; +const DEFAULT_DATA_MOVEMENT_RECHECK_MS: u64 = 250; +const MIN_DATA_MOVEMENT_RECHECK_MS: u64 = 1; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum DataMovementOperation { + Decommission, + Rebalance, +} + +impl DataMovementOperation { + const fn as_str(self) -> &'static str { + match self { + Self::Decommission => "decommission", + Self::Rebalance => "rebalance", + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +struct DataMovementBackpressureConfig { + enabled: bool, + foreground_read_high_percent: usize, + recheck_delay: Duration, +} + +impl DataMovementBackpressureConfig { + fn from_env() -> Self { + Self { + enabled: rustfs_utils::get_env_bool(DATA_MOVEMENT_BACKPRESSURE_ENABLE_ENV, DEFAULT_DATA_MOVEMENT_BACKPRESSURE_ENABLE), + foreground_read_high_percent: rustfs_utils::get_env_usize( + DATA_MOVEMENT_FOREGROUND_READ_HIGH_PERCENT_ENV, + DEFAULT_DATA_MOVEMENT_FOREGROUND_READ_HIGH_PERCENT, + ), + recheck_delay: Duration::from_millis( + rustfs_utils::get_env_u64(DATA_MOVEMENT_RECHECK_MS_ENV, DEFAULT_DATA_MOVEMENT_RECHECK_MS) + .max(MIN_DATA_MOVEMENT_RECHECK_MS), + ), + } + } +} + +pub(crate) async fn wait_for_data_movement_admission( + operation: DataMovementOperation, + pool_index: usize, + cancel_token: &CancellationToken, +) -> Result<()> { + wait_for_data_movement_admission_with_provider( + operation, + pool_index, + cancel_token, + DataMovementBackpressureConfig::from_env(), + runtime_sources::workload_admission_snapshot_provider(), + ) + .await +} + +async fn wait_for_data_movement_admission_with_provider( + operation: DataMovementOperation, + pool_index: usize, + cancel_token: &CancellationToken, + config: DataMovementBackpressureConfig, + provider: Option, +) -> Result<()> { + if cancel_token.is_cancelled() { + return Err(Error::OperationCanceled); + } + if !config.enabled || config.foreground_read_high_percent == 0 { + return Ok(()); + } + let Some(provider) = provider else { + return Ok(()); + }; + + let mut delayed_since: Option = None; + loop { + if cancel_token.is_cancelled() { + record_delay_completion(operation, pool_index, delayed_since, "cancelled"); + return Err(Error::OperationCanceled); + } + + if let Some(usage_pct) = foreground_read_pressure_pct(&config, Some(provider.as_ref())) { + if delayed_since.is_none() { + delayed_since = Some(Instant::now()); + record_delay_start(operation, pool_index, usage_pct, &config); + } + + tokio::select! { + _ = cancel_token.cancelled() => { + record_delay_completion(operation, pool_index, delayed_since, "cancelled"); + return Err(Error::OperationCanceled); + } + _ = sleep(config.recheck_delay) => {} + } + continue; + } + + record_delay_completion(operation, pool_index, delayed_since, "admitted"); + return Ok(()); + } +} + +fn foreground_read_pressure_pct( + config: &DataMovementBackpressureConfig, + provider: Option<&(dyn WorkloadAdmissionSnapshotProvider + Send + Sync)>, +) -> Option { + if !config.enabled || config.foreground_read_high_percent == 0 { + return None; + } + + let snapshot = provider?.workload_admission_snapshot(); + let foreground_read = snapshot.get(WorkloadClass::ForegroundRead)?; + if matches!(foreground_read.state, AdmissionState::Saturated) { + return Some(100); + } + + let limit = foreground_read.limit?; + if limit == 0 { + return None; + } + + let usage_pct = foreground_read + .active + .unwrap_or(0) + .saturating_mul(100) + .checked_div(limit) + .unwrap_or(100); + (usage_pct >= config.foreground_read_high_percent).then_some(usage_pct) +} + +fn record_delay_start( + operation: DataMovementOperation, + pool_index: usize, + foreground_read_usage_pct: usize, + config: &DataMovementBackpressureConfig, +) { + counter!( + "rustfs_data_movement_backpressure_total", + "operation" => operation.as_str().to_string(), + "reason" => "foreground_read_pressure".to_string(), + "result" => "delayed".to_string(), + "pool_index" => pool_index.to_string() + ) + .increment(1); + + debug!( + target: "rustfs::ecstore::data_movement", + event = EVENT_DATA_MOVEMENT_BACKPRESSURE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_DATA_MOVEMENT, + operation = operation.as_str(), + pool_index, + state = "delayed", + reason = "foreground_read_pressure", + foreground_read_usage_pct, + threshold_pct = config.foreground_read_high_percent, + recheck_delay_ms = config.recheck_delay.as_millis(), + "Data movement delayed under foreground read pressure" + ); +} + +fn record_delay_completion( + operation: DataMovementOperation, + pool_index: usize, + delayed_since: Option, + result: &'static str, +) { + let Some(delayed_since) = delayed_since else { + return; + }; + + let delay = delayed_since.elapsed(); + counter!( + "rustfs_data_movement_backpressure_total", + "operation" => operation.as_str().to_string(), + "reason" => "foreground_read_pressure".to_string(), + "result" => result.to_string(), + "pool_index" => pool_index.to_string() + ) + .increment(1); + histogram!( + "rustfs_data_movement_backpressure_delay_seconds", + "operation" => operation.as_str().to_string(), + "reason" => "foreground_read_pressure".to_string(), + "result" => result.to_string(), + "pool_index" => pool_index.to_string() + ) + .record(delay.as_secs_f64()); + + debug!( + target: "rustfs::ecstore::data_movement", + event = EVENT_DATA_MOVEMENT_BACKPRESSURE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_DATA_MOVEMENT, + operation = operation.as_str(), + pool_index, + state = result, + delay_secs = delay.as_secs_f64(), + "Data movement backpressure wait completed" + ); +} + +#[cfg(test)] +mod tests { + use super::*; + use rustfs_concurrency::{WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot}; + use std::sync::Arc; + + #[derive(Debug)] + struct StaticWorkloadProvider { + snapshot: WorkloadAdmissionRegistrySnapshot, + } + + impl StaticWorkloadProvider { + fn new(snapshot: WorkloadAdmissionSnapshot) -> Self { + Self { + snapshot: WorkloadAdmissionRegistrySnapshot::new(vec![snapshot]), + } + } + } + + impl WorkloadAdmissionSnapshotProvider for StaticWorkloadProvider { + fn workload_admission_snapshot(&self) -> WorkloadAdmissionRegistrySnapshot { + self.snapshot.clone() + } + } + + fn test_config() -> DataMovementBackpressureConfig { + DataMovementBackpressureConfig { + enabled: true, + foreground_read_high_percent: 80, + recheck_delay: Duration::from_millis(1), + } + } + + #[test] + fn foreground_read_pressure_treats_saturated_as_full() { + let provider = + StaticWorkloadProvider::new(WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Saturated)); + + assert_eq!(foreground_read_pressure_pct(&test_config(), Some(&provider)), Some(100)); + } + + #[test] + fn foreground_read_pressure_uses_active_limit_threshold() { + let provider = StaticWorkloadProvider::new( + WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Open).with_counts( + Some(8), + None, + Some(10), + ), + ); + + assert_eq!(foreground_read_pressure_pct(&test_config(), Some(&provider)), Some(80)); + } + + #[test] + fn foreground_read_pressure_ignores_open_low_usage() { + let provider = StaticWorkloadProvider::new( + WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Open).with_counts( + Some(7), + None, + Some(10), + ), + ); + + assert_eq!(foreground_read_pressure_pct(&test_config(), Some(&provider)), None); + } + + #[tokio::test] + async fn wait_for_data_movement_admission_returns_when_open() { + let provider: WorkloadSnapshotProviderRef = Arc::new(StaticWorkloadProvider::new( + WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Open).with_counts( + Some(0), + None, + Some(10), + ), + )); + let cancel_token = CancellationToken::new(); + + let result = wait_for_data_movement_admission_with_provider( + DataMovementOperation::Rebalance, + 0, + &cancel_token, + test_config(), + Some(provider), + ) + .await; + + assert!(result.is_ok()); + } + + #[tokio::test] + async fn wait_for_data_movement_admission_cancels_under_pressure() { + let provider: WorkloadSnapshotProviderRef = Arc::new(StaticWorkloadProvider::new(WorkloadAdmissionSnapshot::new( + WorkloadClass::ForegroundRead, + AdmissionState::Saturated, + ))); + let cancel_token = CancellationToken::new(); + let waiter_token = cancel_token.clone(); + + let waiter = tokio::spawn(async move { + wait_for_data_movement_admission_with_provider( + DataMovementOperation::Decommission, + 1, + &waiter_token, + test_config(), + Some(provider), + ) + .await + }); + + tokio::task::yield_now().await; + cancel_token.cancel(); + + let err = waiter + .await + .expect("waiter task should join") + .expect_err("cancelled admission wait should return operation-canceled"); + assert!(matches!(err, Error::OperationCanceled)); + } +} diff --git a/crates/ecstore/src/lib.rs b/crates/ecstore/src/lib.rs index 051fc0ddf..8a16e06fa 100644 --- a/crates/ecstore/src/lib.rs +++ b/crates/ecstore/src/lib.rs @@ -25,6 +25,7 @@ mod cluster; mod compress; mod config; mod data_movement; +mod data_movement_backpressure; mod data_usage; mod disk; mod disks_layout; @@ -61,6 +62,17 @@ mod pools_test; mod store_test; mod tier; +use rustfs_concurrency::WorkloadAdmissionSnapshotProvider; +use std::sync::Arc; + +pub type WorkloadAdmissionSnapshotProviderRef = Arc; + +pub fn set_workload_admission_snapshot_provider( + provider: WorkloadAdmissionSnapshotProviderRef, +) -> std::result::Result<(), WorkloadAdmissionSnapshotProviderRef> { + runtime_sources::set_workload_admission_snapshot_provider(provider) +} + #[cfg(test)] mod rio_tests { #[test] diff --git a/crates/ecstore/src/pools.rs b/crates/ecstore/src/pools.rs index 7c963c980..ae36aedc0 100644 --- a/crates/ecstore/src/pools.rs +++ b/crates/ecstore/src/pools.rs @@ -27,6 +27,7 @@ use crate::bucket::{ use crate::cache_value::metacache_set::{ListPathRawOptions, list_path_raw}; use crate::config::com::{CONFIG_PREFIX, read_config, read_config_no_lock, save_config, save_config_with_opts}; use crate::data_movement; +use crate::data_movement_backpressure::{self, DataMovementOperation}; use crate::data_usage::DATA_USAGE_CACHE_NAME; use crate::disk::error::DiskError; use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}; @@ -3115,7 +3116,40 @@ impl ECStore { let callback_rx = callback_rx.clone(); Box::pin(async move { - let worker_permit = match workers.clone().acquire_owned().await { + if callback_rx.is_cancelled() { + return; + } + if entry_error.lock().await.is_some() { + return; + } + + if let Err(err) = data_movement_backpressure::wait_for_data_movement_admission( + DataMovementOperation::Decommission, + idx, + &callback_rx, + ) + .await + { + if matches!(err, Error::OperationCanceled) { + return; + } + error!("decommission_pool: data movement admission failed: {err}"); + let mut first_err = entry_error.lock().await; + if first_err.is_none() { + *first_err = Some(err); + callback_rx.cancel(); + } + return; + } + + if entry_error.lock().await.is_some() { + return; + } + + let worker_permit = match tokio::select! { + _ = callback_rx.cancelled() => return, + permit = workers.clone().acquire_owned() => permit, + } { Ok(permit) => permit, Err(err) => { let err = Error::other(format!("decommission entry worker permit acquire failed: {err}")); @@ -3128,6 +3162,9 @@ impl ECStore { return; } }; + if entry_error.lock().await.is_some() { + return; + } let entry_rx = callback_rx.clone(); if let Err(err) = this .decommission_entry( diff --git a/crates/ecstore/src/rebalance/entry.rs b/crates/ecstore/src/rebalance/entry.rs index 8a036947a..8545633b5 100644 --- a/crates/ecstore/src/rebalance/entry.rs +++ b/crates/ecstore/src/rebalance/entry.rs @@ -30,6 +30,7 @@ use super::{ REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome, }; use crate::data_movement; +use crate::data_movement_backpressure::{self, DataMovementOperation}; use crate::error::{Error, Result}; use crate::object_api::{GetObjectReader, ObjectOptions}; use crate::pools::ListCallback; @@ -400,6 +401,25 @@ impl ECStore { return; } + if let Err(err) = data_movement_backpressure::wait_for_data_movement_admission( + DataMovementOperation::Rebalance, + pool_index, + &callback_rx, + ) + .await + { + if matches!(err, Error::OperationCanceled) { + return; + } + error!("rebalance_entry: data movement admission failed: {err}"); + let mut first_err = entry_error.lock().await; + if first_err.is_none() { + *first_err = Some(err); + callback_rx.cancel(); + } + return; + } + let permit = tokio::select! { _ = callback_rx.cancelled() => return, permit = entry_workers.clone().acquire_owned() => match permit { diff --git a/crates/ecstore/src/runtime_sources.rs b/crates/ecstore/src/runtime_sources.rs index 3d7c76b19..67eccffe3 100644 --- a/crates/ecstore/src/runtime_sources.rs +++ b/crates/ecstore/src/runtime_sources.rs @@ -12,7 +12,11 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::{collections::HashMap, sync::Arc, time::SystemTime}; +use std::{ + collections::HashMap, + sync::{Arc, OnceLock}, + time::SystemTime, +}; use crate::bucket::bandwidth::monitor::Monitor; use crate::disk::endpoint::Endpoint; @@ -40,6 +44,7 @@ use crate::{ tier::tier::TierConfigMgr, }; use rustfs_common::{GLOBAL_CONN_MAP, GLOBAL_LOCAL_NODE_NAME, GLOBAL_RUSTFS_ADDR, GLOBAL_RUSTFS_HOST}; +use rustfs_concurrency::WorkloadAdmissionSnapshotProvider; use rustfs_config::server_config::{Config, get_global_server_config, set_global_server_config}; use rustfs_io_metrics::internode_metrics::global_internode_metrics; use rustfs_kms::{ObjectEncryptionService, get_global_encryption_service}; @@ -53,6 +58,20 @@ use uuid::Uuid; #[cfg(test)] const TEST_RPC_SECRET: &str = "test-rpc-secret"; +pub(crate) type WorkloadSnapshotProviderRef = Arc; + +static WORKLOAD_ADMISSION_SNAPSHOT_PROVIDER: OnceLock = OnceLock::new(); + +pub(crate) fn set_workload_admission_snapshot_provider( + provider: WorkloadSnapshotProviderRef, +) -> std::result::Result<(), WorkloadSnapshotProviderRef> { + WORKLOAD_ADMISSION_SNAPSHOT_PROVIDER.set(provider) +} + +pub(crate) fn workload_admission_snapshot_provider() -> Option { + WORKLOAD_ADMISSION_SNAPSHOT_PROVIDER.get().cloned() +} + pub(crate) fn record_erasure_write_quorum_failure(stage: &'static str, dominant_error: &'static str) { global_internode_metrics().record_erasure_write_quorum_failure(stage, dominant_error); } diff --git a/rustfs/src/startup_background.rs b/rustfs/src/startup_background.rs index d3ce62b11..f28ec9a80 100644 --- a/rustfs/src/startup_background.rs +++ b/rustfs/src/startup_background.rs @@ -12,8 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::storage_api::startup::ECStore; +use crate::storage_api::startup::{ECStore, set_workload_admission_snapshot_provider}; use crate::workload_admission::RustFsWorkloadAdmissionSnapshotProvider; +use rustfs_concurrency::WorkloadAdmissionSnapshotProvider; use rustfs_heal::{ create_ahm_services_cancel_token, heal::storage::ECStoreHealStorage, init_heal_manager_with_workload_provider, }; @@ -45,9 +46,12 @@ pub(crate) async fn init_background_service_runtime(store: Arc) -> Resu "Background services configured" ); + let workload_provider: Arc = + Arc::new(RustFsWorkloadAdmissionSnapshotProvider); + let _ = set_workload_admission_snapshot_provider(workload_provider.clone()); + if enable_heal || enable_scanner { let heal_storage = Arc::new(ECStoreHealStorage::new(store)); - let workload_provider = Arc::new(RustFsWorkloadAdmissionSnapshotProvider); init_heal_manager_with_workload_provider(heal_storage, None, Some(workload_provider)).await?; } diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 70b547057..08b8c7d59 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -231,6 +231,7 @@ pub(crate) type ObjectLockBlockReason = ecstore_bucket::object_lock::objectlock_ pub(crate) type ObjectStoreResolver = dyn Fn() -> Option> + Send + Sync + 'static; pub(crate) type PolicySys = ecstore_bucket::policy_sys::PolicySys; pub(crate) type PoolEndpoints = ecstore_layout::PoolEndpoints; +pub(crate) type WorkloadAdmissionSnapshotProviderRef = rustfs_ecstore::WorkloadAdmissionSnapshotProviderRef; pub(crate) type QuotaError = ecstore_bucket::quota::QuotaError; pub(crate) type RawFileInfo = rustfs_filemeta::RawFileInfo; pub(crate) type ReadMultipleReq = ecstore_disk::ReadMultipleReq; @@ -850,6 +851,12 @@ pub(crate) fn set_object_store_resolver(resolver: Arc) -> b ecstore_global::set_object_store_resolver(resolver) } +pub(crate) fn set_workload_admission_snapshot_provider( + provider: WorkloadAdmissionSnapshotProviderRef, +) -> std::result::Result<(), WorkloadAdmissionSnapshotProviderRef> { + rustfs_ecstore::set_workload_admission_snapshot_provider(provider) +} + pub(crate) fn get_global_notification_sys() -> Option<&'static NotificationSys> { ecstore_notification::get_global_notification_sys() } diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index d73a6e827..35687cc8d 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -104,8 +104,8 @@ pub(crate) mod startup { init_bucket_metadata_sys, init_ecstore_config, init_global_config_sys, init_local_disks, init_lock_clients, new_global_notification_sys, prewarm_local_disk_id_map, process_lambda_configurations, process_queue_configurations, process_topic_configurations, set_global_endpoints, set_global_region, set_global_rustfs_port, - shutdown_background_services, try_migrate_bucket_metadata, try_migrate_iam_config, try_migrate_server_config, - update_erasure_type, + set_workload_admission_snapshot_provider, shutdown_background_services, try_migrate_bucket_metadata, + try_migrate_iam_config, try_migrate_server_config, update_erasure_type, }; #[cfg(test)]