From 19c7529d88470c207b7bf4cea7f897877122796c Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Fri, 28 Aug 2026 15:11:08 +0800 Subject: [PATCH] refactor(ecstore): migrate data movement backpressure to shared ForegroundPressure (#6779) refactor(ecstore): use shared ForegroundPressure for data movement The data movement backpressure module carried its own byte-identical copy of ForegroundPressure, its reason() label mapping, and the foreground utilization computation. rustfs-concurrency now owns that logic as workload::ForegroundPressure and workload::foreground_pressure, so the local copy was a cross-crate synchronization point that could silently drift from the heal-side and admission-side behavior. Delete the local type and computation and call the shared function instead. The call site keeps what is specific to data movement: the config.enabled short circuit, the optional provider unwrap, and the read/write threshold percentages read from DataMovementBackpressureConfig. The reason() labels emitted into the rustfs_data_movement_backpressure_total metric and the data_movement_backpressure log event are unchanged, as are the existing tests and their assertions. Refs rustfs/backlog#2048 (cherry picked from commit 6a26e144e06ced53a8dfd1712ab7aa24589646ff) (cherry picked from commit ab5ab417e80179265c32b22a5e671eac0b9e43ae) --- .../ecstore/src/data_movement/backpressure.rs | 60 +++---------------- 1 file changed, 8 insertions(+), 52 deletions(-) diff --git a/crates/ecstore/src/data_movement/backpressure.rs b/crates/ecstore/src/data_movement/backpressure.rs index f0406e474..7c9008a5b 100644 --- a/crates/ecstore/src/data_movement/backpressure.rs +++ b/crates/ecstore/src/data_movement/backpressure.rs @@ -15,7 +15,8 @@ use crate::error::{Error, Result}; use crate::runtime::sources::{self as runtime_sources, WorkloadSnapshotProviderRef}; use metrics::{counter, histogram}; -use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshotProvider, WorkloadClass}; +use rustfs_concurrency::WorkloadAdmissionSnapshotProvider; +use rustfs_concurrency::workload::ForegroundPressure; use std::time::{Duration, Instant}; use tokio::time::sleep; use tokio_util::sync::CancellationToken; @@ -137,23 +138,6 @@ async fn wait_for_data_movement_admission_with_provider( } } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -struct ForegroundPressure { - class: WorkloadClass, - usage_pct: usize, - threshold_pct: usize, -} - -impl ForegroundPressure { - const fn reason(self) -> &'static str { - match self.class { - WorkloadClass::ForegroundRead => "foreground_read_pressure", - WorkloadClass::ForegroundWrite => "foreground_write_pressure", - _ => "foreground_pressure", - } - } -} - fn foreground_pressure( config: &DataMovementBackpressureConfig, provider: Option<&(dyn WorkloadAdmissionSnapshotProvider + Send + Sync)>, @@ -163,39 +147,11 @@ fn foreground_pressure( } let snapshot = provider?.workload_admission_snapshot(); - [ - (WorkloadClass::ForegroundRead, config.foreground_read_high_percent), - (WorkloadClass::ForegroundWrite, config.foreground_write_high_percent), - ] - .into_iter() - .filter_map(|(class, threshold_pct)| { - if threshold_pct == 0 { - return None; - } - - let entry = snapshot.get(class)?; - let usage_pct = if matches!(entry.state, AdmissionState::Saturated) { - 100 - } else { - let limit = entry.limit?; - if limit == 0 { - return None; - } - entry - .active - .unwrap_or(0) - .saturating_mul(100) - .checked_div(limit) - .unwrap_or(100) - }; - - (usage_pct >= threshold_pct).then_some(ForegroundPressure { - class, - usage_pct, - threshold_pct, - }) - }) - .max_by_key(|pressure| pressure.usage_pct) + rustfs_concurrency::workload::foreground_pressure( + &snapshot, + config.foreground_read_high_percent, + config.foreground_write_high_percent, + ) } fn record_delay_start( @@ -276,7 +232,7 @@ fn record_delay_completion( #[cfg(test)] mod tests { use super::*; - use rustfs_concurrency::{WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot}; + use rustfs_concurrency::{AdmissionState, WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot, WorkloadClass}; use std::sync::Arc; #[derive(Debug)]