mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-01 17:58:22 +00:00
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)
This commit is contained in:
@@ -15,7 +15,8 @@
|
|||||||
use crate::error::{Error, Result};
|
use crate::error::{Error, Result};
|
||||||
use crate::runtime::sources::{self as runtime_sources, WorkloadSnapshotProviderRef};
|
use crate::runtime::sources::{self as runtime_sources, WorkloadSnapshotProviderRef};
|
||||||
use metrics::{counter, histogram};
|
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 std::time::{Duration, Instant};
|
||||||
use tokio::time::sleep;
|
use tokio::time::sleep;
|
||||||
use tokio_util::sync::CancellationToken;
|
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(
|
fn foreground_pressure(
|
||||||
config: &DataMovementBackpressureConfig,
|
config: &DataMovementBackpressureConfig,
|
||||||
provider: Option<&(dyn WorkloadAdmissionSnapshotProvider + Send + Sync)>,
|
provider: Option<&(dyn WorkloadAdmissionSnapshotProvider + Send + Sync)>,
|
||||||
@@ -163,39 +147,11 @@ fn foreground_pressure(
|
|||||||
}
|
}
|
||||||
|
|
||||||
let snapshot = provider?.workload_admission_snapshot();
|
let snapshot = provider?.workload_admission_snapshot();
|
||||||
[
|
rustfs_concurrency::workload::foreground_pressure(
|
||||||
(WorkloadClass::ForegroundRead, config.foreground_read_high_percent),
|
&snapshot,
|
||||||
(WorkloadClass::ForegroundWrite, config.foreground_write_high_percent),
|
config.foreground_read_high_percent,
|
||||||
]
|
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)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn record_delay_start(
|
fn record_delay_start(
|
||||||
@@ -276,7 +232,7 @@ fn record_delay_completion(
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
use rustfs_concurrency::{WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot};
|
use rustfs_concurrency::{AdmissionState, WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot, WorkloadClass};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
|
|||||||
Reference in New Issue
Block a user