mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
refactor: compose workload admission providers (#3685)
This commit is contained in:
@@ -26,6 +26,8 @@ const NOT_EXPOSED_BY_PROVIDER: &str = "not exposed by RustFS workload admission
|
||||
const REPLICATION_RUNTIME_NOT_INITIALIZED: &str = "replication runtime not initialized";
|
||||
const REPLICATION_QUEUE_STATS_UNAVAILABLE: &str = "replication queue stats unavailable";
|
||||
const SCANNER_ACTIVITY_IDLE_OR_NOT_INITIALIZED: &str = "scanner activity idle or not initialized";
|
||||
const STORAGE_CONCURRENCY_PROVIDER_MISSING_FOREGROUND_READ: &str =
|
||||
"storage concurrency provider did not expose foreground read admission";
|
||||
|
||||
#[derive(Debug, Default, Clone, Copy)]
|
||||
pub struct RustFsWorkloadAdmissionSnapshotProvider;
|
||||
@@ -37,16 +39,25 @@ impl WorkloadAdmissionSnapshotProvider for RustFsWorkloadAdmissionSnapshotProvid
|
||||
}
|
||||
|
||||
pub fn workload_admission_registry_snapshot() -> WorkloadAdmissionRegistrySnapshot {
|
||||
storage_concurrency_workload_admission_snapshot().overlay(runtime_owner_workload_admission_registry_snapshot())
|
||||
}
|
||||
|
||||
fn storage_concurrency_workload_admission_snapshot() -> WorkloadAdmissionRegistrySnapshot {
|
||||
let storage_provider: &dyn WorkloadAdmissionSnapshotProvider = get_concurrency_manager();
|
||||
storage_provider.workload_admission_snapshot()
|
||||
}
|
||||
|
||||
fn runtime_owner_workload_admission_registry_snapshot() -> WorkloadAdmissionRegistrySnapshot {
|
||||
let entries = WorkloadClass::REQUIRED
|
||||
.iter()
|
||||
.copied()
|
||||
.map(|class| match class {
|
||||
WorkloadClass::ForegroundRead => foreground_read_workload_admission_snapshot(),
|
||||
WorkloadClass::Metadata => metadata_workload_admission_snapshot(),
|
||||
WorkloadClass::Scanner => scanner_workload_admission_snapshot(),
|
||||
WorkloadClass::Repair => repair_workload_admission_snapshot(),
|
||||
WorkloadClass::Replication => replication_workload_admission_snapshot(),
|
||||
class => WorkloadAdmissionSnapshot::new(class, AdmissionState::Unknown).with_reason(NOT_EXPOSED_BY_PROVIDER),
|
||||
.filter_map(|class| match class {
|
||||
WorkloadClass::ForegroundRead => None,
|
||||
WorkloadClass::Metadata => Some(metadata_workload_admission_snapshot()),
|
||||
WorkloadClass::Scanner => Some(scanner_workload_admission_snapshot()),
|
||||
WorkloadClass::Repair => Some(repair_workload_admission_snapshot()),
|
||||
WorkloadClass::Replication => Some(replication_workload_admission_snapshot()),
|
||||
class => Some(WorkloadAdmissionSnapshot::new(class, AdmissionState::Unknown).with_reason(NOT_EXPOSED_BY_PROVIDER)),
|
||||
})
|
||||
.collect();
|
||||
|
||||
@@ -54,7 +65,13 @@ pub fn workload_admission_registry_snapshot() -> WorkloadAdmissionRegistrySnapsh
|
||||
}
|
||||
|
||||
pub fn foreground_read_workload_admission_snapshot() -> WorkloadAdmissionSnapshot {
|
||||
get_concurrency_manager().get_object_admission_snapshot()
|
||||
storage_concurrency_workload_admission_snapshot()
|
||||
.get(WorkloadClass::ForegroundRead)
|
||||
.cloned()
|
||||
.unwrap_or_else(|| {
|
||||
WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Unknown)
|
||||
.with_reason(STORAGE_CONCURRENCY_PROVIDER_MISSING_FOREGROUND_READ)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn metadata_workload_admission_snapshot() -> WorkloadAdmissionSnapshot {
|
||||
@@ -192,9 +209,13 @@ mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn foreground_read_snapshot_uses_storage_concurrency_admission() {
|
||||
fn foreground_read_snapshot_uses_storage_concurrency_provider_contract() {
|
||||
let snapshot = foreground_read_workload_admission_snapshot();
|
||||
let storage_snapshot = storage_concurrency_workload_admission_snapshot()
|
||||
.get(WorkloadClass::ForegroundRead)
|
||||
.cloned();
|
||||
|
||||
assert_eq!(Some(snapshot.clone()), storage_snapshot);
|
||||
assert_eq!(snapshot.class, WorkloadClass::ForegroundRead);
|
||||
assert_ne!(snapshot.state, AdmissionState::Unknown);
|
||||
assert!(snapshot.active.is_some());
|
||||
@@ -319,4 +340,23 @@ mod tests {
|
||||
);
|
||||
assert_eq!(registry.get(WorkloadClass::Scanner).and_then(|snapshot| snapshot.active), Some(0));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn provider_overlays_runtime_owner_snapshots_on_storage_registry() {
|
||||
let storage_registry = storage_concurrency_workload_admission_snapshot();
|
||||
let registry = workload_admission_registry_snapshot();
|
||||
|
||||
assert_eq!(
|
||||
registry.get(WorkloadClass::ForegroundRead),
|
||||
storage_registry.get(WorkloadClass::ForegroundRead)
|
||||
);
|
||||
assert_eq!(registry.get(WorkloadClass::Metadata), Some(&metadata_workload_admission_snapshot()));
|
||||
assert_eq!(registry.get(WorkloadClass::Scanner), Some(&scanner_workload_admission_snapshot()));
|
||||
assert_eq!(
|
||||
registry
|
||||
.get(WorkloadClass::ForegroundWrite)
|
||||
.and_then(|snapshot| snapshot.reason.as_deref()),
|
||||
Some(NOT_EXPOSED_BY_PROVIDER)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user