feat(scanner): coordinate usage and workload boundaries (#7093)

* test(scanner): wire usage and heal rebuild gates

* docs(scanner): define usage authority protocol

* docs(heal): clarify scanner and ecstore boundaries

* refactor(scanner): split metrics from contracts

* feat(scanner): use shared workload snapshots

* fix(ecstore): recheck capacity before decommission drain
This commit is contained in:
houseme
2026-09-03 17:02:43 +08:00
committed by GitHub
parent 3ab7a1921f
commit 0e6ee3bf62
57 changed files with 632 additions and 134 deletions
+1 -1
View File
@@ -130,7 +130,7 @@ rustfs-concurrency.workspace = true
rustfs-credentials = { workspace = true }
rustfs-common.workspace = true
rustfs-heal-contracts.workspace = true
rustfs-scanner-contracts.workspace = true
rustfs-scanner-metrics.workspace = true
rustfs-policy.workspace = true
rustfs-protos.workspace = true
rustfs-replication.workspace = true
@@ -17,7 +17,7 @@ use crate::bucket::lifecycle::lifecycle;
use crate::object_api::ObjectInfo;
use crate::services::event_notification::{EventArgs, send_event};
use rustfs_s3_types::EventName;
use rustfs_scanner_contracts::metrics::IlmAction;
use rustfs_scanner_metrics::metrics::IlmAction;
const LIFECYCLE_EXPIRY_USER_AGENT: &str = "Internal: [ILM-Expiry]";
const LIFECYCLE_TRANSITION_USER_AGENT: &str = "Internal: [ILM-Transition]";
@@ -84,7 +84,7 @@ use rustfs_data_usage::TierStats;
use rustfs_filemeta::{
FileInfo, FileInfoOpts, NULL_VERSION_ID, RestoreStatusOps, TRANSITION_COMPLETE, get_file_info, is_restored_object_on_disk,
};
use rustfs_scanner_contracts::metrics::{
use rustfs_scanner_metrics::metrics::{
IlmAction, Metrics, ScannerLifecycleExpiryStateUpdate, ScannerLifecycleTransitionStateUpdate, global_metrics,
};
use rustfs_utils::{
@@ -5596,7 +5596,7 @@ mod tests {
use rustfs_filemeta::{FileInfo, FileMeta};
#[cfg(feature = "test-util")]
use rustfs_s3_client::transition_api::ReaderImpl;
use rustfs_scanner_contracts::metrics::{IlmAction, global_metrics};
use rustfs_scanner_metrics::metrics::{IlmAction, global_metrics};
use s3s::dto::{
BucketLifecycleConfiguration, DefaultRetention, ExpirationStatus, LifecycleExpiration, LifecycleRule, MetadataEntry,
ObjectLockConfiguration, ObjectLockEnabled, ObjectLockRetentionMode, ObjectLockRule, OutputLocation, RestoreRequest,
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use rustfs_scanner_contracts::metrics::IlmAction;
use rustfs_scanner_metrics::metrics::IlmAction;
use crate::bucket::lifecycle::lifecycle::ObjectOpts;
use crate::bucket::replication::ReplicationLifecycleBridge;
@@ -77,7 +77,7 @@ mod tests {
use crate::bucket::replication::{DeleteReplicationConfigSnapshot, ReplicationObjectBridge};
use crate::object_api::{ObjectInfo, ObjectOptions};
use crate::storage_api_contracts::object::ObjectToDelete;
use rustfs_scanner_contracts::metrics::IlmAction;
use rustfs_scanner_metrics::metrics::IlmAction;
use s3s::dto::{
BucketVersioningStatus, DeleteMarkerReplication, DeleteMarkerReplicationStatus, DeleteReplication,
DeleteReplicationStatus, Destination, ReplicationConfiguration, ReplicationRule, ReplicationRuleStatus,
+5 -5
View File
@@ -16,7 +16,7 @@ use super::{BucketQuota, QuotaCheckResult, QuotaError, QuotaOperation};
use crate::bucket::metadata_sys::{BucketMetadataSys, update, update_if_incarnation};
use crate::data_usage::get_bucket_usage_memory;
use rustfs_config::QUOTA_CONFIG_FILE;
use rustfs_scanner_contracts::metrics::Metric;
use rustfs_scanner_metrics::metrics::Metric;
use std::sync::Arc;
use std::time::Instant;
use time::OffsetDateTime;
@@ -120,9 +120,9 @@ impl QuotaChecker {
let duration = start_time.elapsed();
// inc_time is now a plain fn (not async) — no .await needed.
rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaCheck, duration);
rustfs_scanner_metrics::metrics::Metrics::inc_time(Metric::QuotaCheck, duration);
if !allowed {
rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaViolation, duration);
rustfs_scanner_metrics::metrics::Metrics::inc_time(Metric::QuotaViolation, duration);
}
Ok(result)
@@ -185,7 +185,7 @@ impl QuotaChecker {
.await
.map_err(QuotaError::StorageError)?;
rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed());
rustfs_scanner_metrics::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed());
Ok(updated_at)
}
@@ -206,7 +206,7 @@ impl QuotaChecker {
}
.map_err(QuotaError::StorageError)?;
rustfs_scanner_contracts::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed());
rustfs_scanner_metrics::metrics::Metrics::inc_time(Metric::QuotaSync, start_time.elapsed());
Ok(updated_at)
}
+43 -1
View File
@@ -13504,7 +13504,9 @@ impl ECStore {
}
return Ok(());
}
let result = self.decommission_in_background(rx.clone(), idx, entry_budget).await;
let result = self
.decommission_in_background(rx.clone(), idx, generation, entry_budget)
.await;
if let Err(err) = &result
&& (is_decommission_capacity_blocked_error(err) || is_decommission_target_capacity_error(err))
@@ -13989,8 +13991,10 @@ impl ECStore {
self: &Arc<Self>,
rx: CancellationToken,
idx: usize,
generation: OffsetDateTime,
entry_budget: Arc<Semaphore>,
) -> Result<()> {
self.ensure_decommission_runtime_capacity_available(idx, generation).await?;
let pool = get_by_index(self.pools.as_slice(), idx, "load decommission background pool")?.clone();
let pending = {
@@ -17253,6 +17257,44 @@ mod tests {
);
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_worker_rechecks_runtime_capacity_before_empty_background_completion() {
let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await;
let layout = DecommissionErasureLayout { data: 1, parity: 0 };
let enough = vec![
DecommissionPoolCapacityInfo::for_test(0, layout, 0, 30, 30),
DecommissionPoolCapacityInfo::for_test(1, layout, 60, 60, 0),
];
set_decommission_capacity_info_overrides_for_test(store.id, vec![enough.clone()]);
store
.save_current_pool_meta_for_decommission_start(&[0], Vec::new())
.await
.expect("the initial reservation should be activated");
let shortage = vec![enough[0], DecommissionPoolCapacityInfo::for_test(1, layout, 59, 60, 1)];
set_decommission_capacity_info_overrides_for_test(store.id, vec![shortage]);
let canceler = DecommissionCanceler::new(CancellationToken::new());
store.decommission_cancelers.write().await[0] = Some(canceler.clone());
store
.do_decommission_in_routine(canceler, 0, Arc::new(Semaphore::new(1)))
.await
.expect("runtime capacity shortage should pause the worker before background completion");
let local = store.pool_meta.read().await;
let info = local.pools[0]
.decommission
.as_ref()
.expect("the blocked decommission state should remain present");
assert!(!info.complete && !info.failed && !info.canceled);
assert!(info.capacity_blocked_reason.is_some());
assert!(
info.capacity_reservation
.as_ref()
.is_some_and(DecommissionCapacityReservation::active)
);
}
fn pool_meta_replica_test_meta(cmd_line: &str) -> PoolMeta {
PoolMeta {
version: POOL_META_VERSION,
+1 -1
View File
@@ -286,7 +286,7 @@ pub struct QuotaAdmission {
pub struct LifecycleDeleteAllRequest {
pub(crate) version_id: Option<Uuid>,
pub(crate) delete_marker: bool,
pub(crate) action: rustfs_scanner_contracts::metrics::IlmAction,
pub(crate) action: rustfs_scanner_metrics::metrics::IlmAction,
pub(crate) rule_id: String,
pub(crate) phase: LifecycleDeleteAllPhase,
}
+26 -26
View File
@@ -32,7 +32,7 @@ use rustfs_madmin::metrics::{
ScannerSourceCycleSnapshot as MadminScannerSourceCycleSnapshot, ScannerSourceWorkSnapshot as MadminScannerSourceWorkSnapshot,
ScannerUsageFreshnessSnapshot as MadminScannerUsageFreshnessSnapshot, TimedAction as MadminTimedAction,
};
use rustfs_scanner_contracts::metrics::global_metrics;
use rustfs_scanner_metrics::metrics::global_metrics;
use rustfs_utils::os::get_drive_stats;
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
@@ -82,7 +82,7 @@ fn unix_millis_to_jiff_timestamp(millis: u64, fallback: Timestamp) -> Timestamp
}
}
fn to_madmin_scanner_metrics(metrics: rustfs_scanner_contracts::metrics::ScannerMetricsReport) -> MadminScannerMetrics {
fn to_madmin_scanner_metrics(metrics: rustfs_scanner_metrics::metrics::ScannerMetricsReport) -> MadminScannerMetrics {
MadminScannerMetrics {
collected_at: metrics.collected_at,
current_cycle: metrics.current_cycle,
@@ -565,7 +565,7 @@ async fn collect_local_disks_metrics(disks: &HashSet<String>) -> HashMap<String,
mod test {
use super::*;
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_scanner_contracts::metrics::CurrentCycle;
use rustfs_scanner_metrics::metrics::CurrentCycle;
use serial_test::serial;
use std::time::Duration;
@@ -619,7 +619,7 @@ mod test {
#[test]
fn scanner_metrics_mapping_preserves_partial_source_status() {
let current_started = Utc::now() - chrono::Duration::seconds(5);
let scanner = to_madmin_scanner_metrics(rustfs_scanner_contracts::metrics::ScannerMetricsReport {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_metrics::metrics::ScannerMetricsReport {
current_cycle_active: true,
current_started: chrono_to_jiff_timestamp(current_started),
last_cycle_partial_source: "usage".to_string(),
@@ -628,7 +628,7 @@ mod test {
cycle_recovery_required_total: 2,
cycle_last_progress_age: 17,
leader_lease_without_progress: true,
partial_cycles_by_source: vec![rustfs_scanner_contracts::metrics::ScannerSourceCycleSnapshot {
partial_cycles_by_source: vec![rustfs_scanner_metrics::metrics::ScannerSourceCycleSnapshot {
source: "usage".to_string(),
cycles: 2,
}],
@@ -654,11 +654,11 @@ mod test {
#[tokio::test]
#[serial]
async fn collect_local_metrics_preserves_scanner_cycle_started_time() {
let previous_init_time = *rustfs_scanner_contracts::GLOBAL_INIT_TIME.read().await;
let previous_init_time = *rustfs_scanner_metrics::GLOBAL_INIT_TIME.read().await;
let previous_cycle = global_metrics().get_cycle().await;
let init_time = Utc::now() - chrono::Duration::hours(1);
let cycle_started = Utc::now() - chrono::Duration::seconds(5);
*rustfs_scanner_contracts::GLOBAL_INIT_TIME.write().await = Some(init_time);
*rustfs_scanner_metrics::GLOBAL_INIT_TIME.write().await = Some(init_time);
let cycle = CurrentCycle {
current: 0,
next: 1,
@@ -673,7 +673,7 @@ mod test {
.finish_scan_cycle_work_with_cycle(cycle_start, previous_cycle.clone().unwrap_or_default())
.await;
global_metrics().set_cycle(previous_cycle).await;
*rustfs_scanner_contracts::GLOBAL_INIT_TIME.write().await = previous_init_time;
*rustfs_scanner_metrics::GLOBAL_INIT_TIME.write().await = previous_init_time;
let encoded = rmp_serde::to_vec_named(&realtime).expect("realtime metrics should encode");
let decoded: RealtimeMetrics = rmp_serde::from_slice(&encoded).expect("realtime metrics should decode");
@@ -686,8 +686,8 @@ mod test {
#[test]
fn scanner_metrics_mapping_preserves_pacing_pressure() {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_contracts::metrics::ScannerMetricsReport {
pacing_pressure: rustfs_scanner_contracts::metrics::ScannerPacingPressureSnapshot {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_metrics::metrics::ScannerMetricsReport {
pacing_pressure: rustfs_scanner_metrics::metrics::ScannerPacingPressureSnapshot {
primary_pressure: "cycle_budget".to_string(),
current_queued_scans: 4,
current_active_scans: 2,
@@ -712,12 +712,12 @@ mod test {
#[test]
fn scanner_metrics_mapping_preserves_lifecycle_transition_status() {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_contracts::metrics::ScannerMetricsReport {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_metrics::metrics::ScannerMetricsReport {
current_cycle_lifecycle_expiry_actions: 2,
current_cycle_lifecycle_transition_actions: 3,
last_cycle_lifecycle_expiry_actions: 5,
last_cycle_lifecycle_transition_actions: 7,
lifecycle_expiry: rustfs_scanner_contracts::metrics::ScannerLifecycleExpirySnapshot {
lifecycle_expiry: rustfs_scanner_metrics::metrics::ScannerLifecycleExpirySnapshot {
current_queue_capacity: 16,
current_queued: 5,
current_active: 2,
@@ -729,7 +729,7 @@ mod test {
scanner_not_enqueued: 2,
delete_failed: 1,
},
lifecycle_transition: rustfs_scanner_contracts::metrics::ScannerLifecycleTransitionSnapshot {
lifecycle_transition: rustfs_scanner_metrics::metrics::ScannerLifecycleTransitionSnapshot {
current_queue_capacity: 16,
current_queued: 5,
current_active: 2,
@@ -778,10 +778,10 @@ mod test {
#[test]
fn scanner_metrics_mapping_preserves_maintenance_control_status() {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_contracts::metrics::ScannerMetricsReport {
maintenance_control: rustfs_scanner_contracts::metrics::ScannerMaintenanceControlSnapshot {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_metrics::metrics::ScannerMetricsReport {
maintenance_control: rustfs_scanner_metrics::metrics::ScannerMaintenanceControlSnapshot {
primary_control: "blocked_source".to_string(),
sources: vec![rustfs_scanner_contracts::metrics::ScannerMaintenanceSourceSnapshot {
sources: vec![rustfs_scanner_metrics::metrics::ScannerMaintenanceSourceSnapshot {
source: "lifecycle".to_string(),
state: "blocked".to_string(),
reason: "missed_work".to_string(),
@@ -815,8 +815,8 @@ mod test {
#[test]
fn scanner_metrics_mapping_preserves_usage_freshness_status() {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_contracts::metrics::ScannerMetricsReport {
usage_freshness: rustfs_scanner_contracts::metrics::ScannerUsageFreshnessSnapshot {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_metrics::metrics::ScannerMetricsReport {
usage_freshness: rustfs_scanner_metrics::metrics::ScannerUsageFreshnessSnapshot {
dirty_pending_buckets: 3,
last_dirty_mark_unix_secs: 10,
last_dirty_clear_unix_secs: 11,
@@ -857,7 +857,7 @@ mod test {
#[test]
fn scanner_metrics_mapping_preserves_distributed_status_fields() {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_contracts::metrics::ScannerMetricsReport {
let scanner = to_madmin_scanner_metrics(rustfs_scanner_metrics::metrics::ScannerMetricsReport {
active_scan_paths: 2,
oldest_active_path_age_seconds: 45,
active_paths: vec!["disk-a/bucket-a".to_string(), "disk-b/bucket-b".to_string()],
@@ -913,7 +913,7 @@ mod test {
cycle_max_directories: 38,
bitrot_cycle_enabled: true,
bitrot_cycle_seconds: 39.0,
scan_checkpoint: Some(rustfs_scanner_contracts::metrics::ScannerCheckpointReport {
scan_checkpoint: Some(rustfs_scanner_metrics::metrics::ScannerCheckpointReport {
version: 1,
resume_after: "bucket-a/prefix-a".to_string(),
reason: "directories".to_string(),
@@ -923,7 +923,7 @@ mod test {
scan_checkpoint_cleared: 41,
scan_checkpoint_ignored: 42,
scan_checkpoint_stale: 43,
source_work: vec![rustfs_scanner_contracts::metrics::ScannerSourceWorkSnapshot {
source_work: vec![rustfs_scanner_metrics::metrics::ScannerSourceWorkSnapshot {
source: "usage".to_string(),
checked: 44,
queued: 45,
@@ -932,7 +932,7 @@ mod test {
skipped: 48,
missed: 49,
}],
current_cycle_source_work: vec![rustfs_scanner_contracts::metrics::ScannerSourceWorkSnapshot {
current_cycle_source_work: vec![rustfs_scanner_metrics::metrics::ScannerSourceWorkSnapshot {
source: "lifecycle".to_string(),
checked: 50,
queued: 51,
@@ -941,7 +941,7 @@ mod test {
skipped: 54,
missed: 55,
}],
last_cycle_source_work: vec![rustfs_scanner_contracts::metrics::ScannerSourceWorkSnapshot {
last_cycle_source_work: vec![rustfs_scanner_metrics::metrics::ScannerSourceWorkSnapshot {
source: "heal".to_string(),
checked: 56,
queued: 57,
@@ -950,7 +950,7 @@ mod test {
skipped: 60,
missed: 61,
}],
replication_repair: vec![rustfs_scanner_contracts::metrics::ScannerReplicationRepairSnapshot {
replication_repair: vec![rustfs_scanner_metrics::metrics::ScannerReplicationRepairSnapshot {
source: "bucket_replication".to_string(),
kind: "object".to_string(),
scanner_role: "repair_admission".to_string(),
@@ -962,7 +962,7 @@ mod test {
skipped: 66,
missed: 67,
}],
current_cycle_replication_repair: vec![rustfs_scanner_contracts::metrics::ScannerReplicationRepairSnapshot {
current_cycle_replication_repair: vec![rustfs_scanner_metrics::metrics::ScannerReplicationRepairSnapshot {
source: "bucket_replication".to_string(),
kind: "delete_marker".to_string(),
scanner_role: "repair_admission".to_string(),
@@ -974,7 +974,7 @@ mod test {
skipped: 72,
missed: 73,
}],
last_cycle_replication_repair: vec![rustfs_scanner_contracts::metrics::ScannerReplicationRepairSnapshot {
last_cycle_replication_repair: vec![rustfs_scanner_metrics::metrics::ScannerReplicationRepairSnapshot {
source: "site_replication".to_string(),
kind: "active_resync".to_string(),
scanner_role: "boundary_signal".to_string(),
+3 -3
View File
@@ -1168,7 +1168,7 @@ mod lifecycle_delete_all_plan_tests {
crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(version_id),
delete_marker: true,
action: rustfs_scanner_contracts::metrics::IlmAction::DelMarkerDeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DelMarkerDeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}
@@ -1361,7 +1361,7 @@ mod lifecycle_delete_all_plan_tests {
let request = crate::object_api::LifecycleDeleteAllRequest {
version_id: None,
delete_marker: false,
action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
};
@@ -19363,7 +19363,7 @@ mod delete_objects_lock_gating_tests {
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(trigger_version_id),
delete_marker: false,
action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::History,
}),
+2 -2
View File
@@ -15934,7 +15934,7 @@ mod tests {
.expect("unknown transition metadata should be written");
}
let lifecycle_event = crate::bucket::lifecycle::lifecycle::Event {
action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "delete-all-versions".to_string(),
..Default::default()
};
@@ -16186,7 +16186,7 @@ mod tests {
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: original.version_id,
delete_marker: false,
action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
+1 -1
View File
@@ -6448,7 +6448,7 @@ mod tests {
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(version_id),
delete_marker: false,
action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "delete-all".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
+2 -2
View File
@@ -1322,7 +1322,7 @@ mod tests {
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(trigger_id),
delete_marker: false,
action: rustfs_scanner_contracts::metrics::IlmAction::DeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
@@ -1492,7 +1492,7 @@ mod tests {
lifecycle_delete_all: Some(crate::object_api::LifecycleDeleteAllRequest {
version_id: Some(marker_id),
delete_marker: true,
action: rustfs_scanner_contracts::metrics::IlmAction::DelMarkerDeleteAllVersionsAction,
action: rustfs_scanner_metrics::metrics::IlmAction::DelMarkerDeleteAllVersionsAction,
rule_id: "rule".to_string(),
phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
}),
+1 -1
View File
@@ -53,7 +53,7 @@ hotpath-cpu = [
hotpath.workspace = true
async-trait.workspace = true
metrics.workspace = true
rustfs-scanner-contracts.workspace = true
rustfs-scanner-metrics.workspace = true
rustfs-config = { workspace = true, features = ["constants"] }
rustfs-replication.workspace = true
rustfs-storage-api.workspace = true
+1 -1
View File
@@ -62,7 +62,7 @@ const ERR_LIFECYCLE_EXPIRED_OBJECT_DELETE_MARKER_WITH_TAGS: &str =
const ERR_LIFECYCLE_RULE_MUST_HAVE_ACTION: &str = "Rule must have at least one of Expiration, Transition, NoncurrentVersionExpiration, NoncurrentVersionTransition, or DelMarkerExpiration";
const ERR_LIFECYCLE_PREFIX_FILTER_CONFLICT: &str = "Legacy Prefix and Filter cannot both be present in a lifecycle rule. Use Filter.Prefix instead of the top-level Prefix element.";
pub use rustfs_scanner_contracts::metrics::IlmAction;
pub use rustfs_scanner_metrics::metrics::IlmAction;
#[async_trait::async_trait]
pub trait RuleValidate {
+2 -2
View File
@@ -19,7 +19,7 @@ use time::OffsetDateTime;
use tracing::info;
use rustfs_replication::ReplicationStatusType;
use rustfs_scanner_contracts::metrics::IlmAction;
use rustfs_scanner_metrics::metrics::IlmAction;
use crate::object_lock;
use crate::{Event, Lifecycle, ObjectOpts};
@@ -197,7 +197,7 @@ mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use rustfs_scanner_contracts::metrics::IlmAction;
use rustfs_scanner_metrics::metrics::IlmAction;
use s3s::dto::{
BucketLifecycleConfiguration, DefaultRetention, ExpirationStatus, LifecycleExpiration, LifecycleRule,
NoncurrentVersionExpiration, ObjectLockConfiguration, ObjectLockEnabled, ObjectLockRetentionMode, ObjectLockRule,
+1 -1
View File
@@ -21,4 +21,4 @@ mod tagging;
pub use core::*;
pub use evaluator::Evaluator;
pub use rustfs_replication::{ReplicationStatusType, VersionPurgeStatusType};
pub use rustfs_scanner_contracts::metrics::IlmAction;
pub use rustfs_scanner_metrics::metrics::IlmAction;
+1 -1
View File
@@ -106,7 +106,7 @@ hotpath.workspace = true
rustfs-audit = { workspace = true }
rustfs-common = { workspace = true }
rustfs-heal-contracts = { workspace = true }
rustfs-scanner-contracts = { workspace = true }
rustfs-scanner-metrics = { workspace = true }
rustfs-config = { workspace = true, features = ["observability"] }
# NOTE: This dependency on rustfs-ecstore is a known architectural limitation.
# The obs crate imports types from ecstore for metrics collection.
+1 -1
View File
@@ -493,7 +493,7 @@ mod tests {
use super::*;
use crate::metrics::report::report_metrics;
use metrics_util::debugging::DebuggingRecorder;
use rustfs_scanner_contracts::metrics::{Metric, Metrics};
use rustfs_scanner_metrics::metrics::{Metric, Metrics};
fn prometheus_counter_name(name: &str) -> String {
if name.ends_with("_total") {
+4 -4
View File
@@ -45,7 +45,7 @@ use rustfs_io_metrics::{
ProcessResourceSnapshot, ProcessSampler, ProcessStatusSnapshot, ProcessSystemSnapshot, s3_op_metrics_snapshot,
snapshot_process_resource_and_system, snapshot_process_resource_and_system_with,
};
use rustfs_scanner_contracts::metrics::{
use rustfs_scanner_metrics::metrics::{
ScannerActiveBucketDriveSnapshot, ScannerBucketDriveResultSnapshot, ScannerMetricsReport, ScannerSourceWorkSnapshot,
global_metrics,
};
@@ -1679,7 +1679,7 @@ pub async fn collect_compression_cluster_stats() -> Option<CompressionClusterSta
#[cfg(test)]
mod tests {
use super::*;
use rustfs_scanner_contracts::metrics::ScannerSourceWorkSnapshot;
use rustfs_scanner_metrics::metrics::ScannerSourceWorkSnapshot;
use std::io::{Read, Write};
use std::net::{Shutdown, TcpListener, TcpStream};
use std::thread;
@@ -2134,7 +2134,7 @@ mod tests {
#[test]
fn ilm_detail_stats_keep_expiry_and_transition_results_separate() {
let report = ScannerMetricsReport {
lifecycle_expiry: rustfs_scanner_contracts::metrics::ScannerLifecycleExpirySnapshot {
lifecycle_expiry: rustfs_scanner_metrics::metrics::ScannerLifecycleExpirySnapshot {
current_queued: 2,
current_active: 1,
scanner_queued: 10,
@@ -2142,7 +2142,7 @@ mod tests {
delete_failed: 4,
..Default::default()
},
lifecycle_transition: rustfs_scanner_contracts::metrics::ScannerLifecycleTransitionSnapshot {
lifecycle_transition: rustfs_scanner_metrics::metrics::ScannerLifecycleTransitionSnapshot {
current_queued: 5,
current_active: 6,
queue_full: 7,
+2 -12
View File
@@ -20,26 +20,16 @@ license.workspace = true
repository.workspace = true
rust-version.workspace = true
homepage.workspace = true
description = "Scanner metrics and lifecycle-cycle contracts shared by the scanner, storage engine, and observability layers."
keywords = ["scanner", "contracts", "metrics", "rustfs", "Minio"]
description = "Scanner storage and wire contracts shared by the scanner and storage engine."
keywords = ["scanner", "contracts", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "data-structures"]
[lints]
workspace = true
[dependencies]
chrono = { workspace = true, features = ["serde"] }
jiff = { workspace = true, features = ["serde"] }
metrics = { workspace = true }
rmp-serde = { workspace = true }
rustfs-heal-contracts = { workspace = true }
serde = { workspace = true, features = ["derive"] }
tokio = { workspace = true, features = ["sync"] }
[dev-dependencies]
serde_json = { workspace = true }
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }
uuid = { workspace = true, features = ["v4"] }
[lib]
doctest = false
-6
View File
@@ -11,9 +11,3 @@
// 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.
pub mod init_time;
pub mod last_minute;
pub mod metrics;
pub use init_time::{GLOBAL_INIT_TIME, get_global_init_time, set_global_init_time_now};
+45
View File
@@ -0,0 +1,45 @@
# 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.
[package]
name = "rustfs-scanner-metrics"
version.workspace = true
edition.workspace = true
license.workspace = true
repository.workspace = true
rust-version.workspace = true
homepage.workspace = true
description = "Scanner metrics and lifecycle-cycle telemetry shared by the scanner, storage engine, and observability layers."
keywords = ["scanner", "metrics", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "data-structures"]
[lints]
workspace = true
[dependencies]
chrono = { workspace = true, features = ["serde"] }
jiff = { workspace = true, features = ["serde"] }
metrics = { workspace = true }
rmp-serde = { workspace = true }
rustfs-heal-contracts = { workspace = true }
serde = { workspace = true, features = ["derive"] }
tokio = { workspace = true, features = ["sync"] }
[dev-dependencies]
serde_json = { workspace = true }
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }
uuid = { workspace = true, features = ["v4"] }
[lib]
doctest = false
+19
View File
@@ -0,0 +1,19 @@
// 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.
pub mod init_time;
pub mod last_minute;
pub mod metrics;
pub use init_time::{GLOBAL_INIT_TIME, get_global_init_time, set_global_init_time_now};
+2 -1
View File
@@ -73,8 +73,9 @@ hotpath-cpu = [
hotpath.workspace = true
rustfs-config = { workspace = true, features = ["server-config-model"] }
rustfs-common = { workspace = true }
rustfs-concurrency = { workspace = true }
rustfs-heal-contracts = { workspace = true }
rustfs-scanner-contracts = { workspace = true }
rustfs-scanner-metrics = { workspace = true }
rustfs-credentials = { workspace = true }
rustfs-utils = { workspace = true }
tokio = { workspace = true, features = ["fs", "sync", "time", "macros", "rt-multi-thread"] }
+3 -1
View File
@@ -70,6 +70,7 @@ mod scanner_heal_admission_baseline;
pub mod scanner_io;
pub mod sleeper;
pub(crate) mod storage_api;
mod workload_admission;
pub use data_usage_define::*;
pub use error::ScannerError;
@@ -80,7 +81,7 @@ pub use remote_scanner::{
remote_scanner_request_matches_envelope, serve_remote_scanner_request, validate_remote_scanner_request_fence,
};
pub use runtime_config::{apply_scanner_runtime_config, scanner_runtime_config_status, validate_scanner_runtime_config};
pub use rustfs_scanner_contracts::last_minute;
pub use rustfs_scanner_metrics::last_minute;
pub use scanner::{
ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerCycleScheduleStatus, ScannerPauseBacklogAlertReason,
ScannerPauseBacklogPhase, ScannerPauseBacklogStatus, ScannerPauseBacklogThresholds, ScannerUsageStateResetResult,
@@ -96,6 +97,7 @@ pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER};
use std::sync::atomic::{AtomicU64, Ordering};
pub use storage_api::ScannerReplicationConfig as ReplicationConfig;
pub use storage_api::scan::{SCANNER_ACTIVITY_PROTOCOL_VERSION, SCANNER_ACTIVITY_V6_PROTOCOL_VERSION};
pub use workload_admission::set_scanner_workload_admission_snapshot_provider;
static SCANNER_ACTIVE_WORK_UNITS: AtomicU64 = AtomicU64::new(0);
static SCANNER_RUNTIME_INSTANCES: AtomicU64 = AtomicU64::new(0);
+1 -1
View File
@@ -28,7 +28,7 @@ use crate::{
use hmac::{Hmac, KeyInit, Mac};
use rustfs_credentials::try_get_rpc_token;
use rustfs_heal_contracts::heal_channel::HealScanMode;
use rustfs_scanner_contracts::metrics::{Metric, Metrics};
use rustfs_scanner_metrics::metrics::{Metric, Metrics};
use rustfs_utils::path::path_join_buf;
use serde::{Deserialize, Serialize};
use sha2::Sha256;
+1 -1
View File
@@ -52,7 +52,7 @@ use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELA
use rustfs_data_usage::observed_data_usage_is_newer;
use rustfs_heal_contracts::heal_channel::HealScanMode;
use rustfs_lock::{NamespaceLockGuard, error::LockError};
use rustfs_scanner_contracts::metrics::{
use rustfs_scanner_metrics::metrics::{
CurrentCycle, Metric, Metrics, ScanCyclePartialReason, ScanCycleWorkSnapshot, ScannerUsageSaveResult, ScannerWorkSource,
emit_scan_cycle_complete, emit_scan_cycle_deferred, emit_scan_cycle_partial_with_source, emit_scan_cycle_superseded,
global_metrics,
+1 -1
View File
@@ -50,7 +50,7 @@ use rustfs_heal_contracts::heal_channel::{
HEAL_DELETE_DANGLING, HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest,
HealRequestSource, HealScanMode, send_heal_request_with_admission,
};
use rustfs_scanner_contracts::metrics::{
use rustfs_scanner_metrics::metrics::{
CloseDiskGuard, IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource,
UpdateCurrentPathFn, current_path_updater, global_metrics,
};
+1 -1
View File
@@ -30,7 +30,7 @@ use rustfs_data_usage::{BucketTargetUsageInfo, BucketUsageInfo};
use rustfs_filemeta::FileMeta;
use rustfs_heal_contracts::heal_channel::HealScanMode;
use rustfs_lock::{LockError, NamespaceLockGuard};
use rustfs_scanner_contracts::metrics::{
use rustfs_scanner_metrics::metrics::{
Metric, Metrics, emit_scan_bucket_drive_complete, emit_scan_bucket_drive_partial, global_metrics,
};
use rustfs_utils::path::path_join_buf;
+4 -4
View File
@@ -147,13 +147,13 @@ impl Drop for DiskBucketScanActiveGuard {
pub(super) struct BucketDriveFailureGuard {
failed: bool,
source: rustfs_scanner_contracts::metrics::ScannerWorkSource,
source: rustfs_scanner_metrics::metrics::ScannerWorkSource,
bucket: String,
drive: String,
}
impl BucketDriveFailureGuard {
pub(super) fn new(source: rustfs_scanner_contracts::metrics::ScannerWorkSource, bucket: &str, drive: &str) -> Self {
pub(super) fn new(source: rustfs_scanner_metrics::metrics::ScannerWorkSource, bucket: &str, drive: &str) -> Self {
Self {
failed: true,
source,
@@ -245,7 +245,7 @@ pub(super) fn scanner_concurrency_limit(configured: usize, available: usize) ->
return 0;
}
if crate::current_foreground_read_activity() > 0 {
if crate::workload_admission::foreground_workload_activity() > 0 {
return 1;
}
@@ -285,7 +285,7 @@ pub(super) fn scanner_task_join_error(stage: &str, err: tokio::task::JoinError)
#[cfg(test)]
mod tests {
use super::*;
use rustfs_scanner_contracts::metrics::{ScannerWorkSource, global_metrics};
use rustfs_scanner_metrics::metrics::{ScannerWorkSource, global_metrics};
use tokio::sync::oneshot;
fn active_bucket_drive_count(source: ScannerWorkSource, bucket: &str, drive: &str) -> u64 {
+2 -2
View File
@@ -161,8 +161,8 @@ impl ScannerIODisk for Disk {
let bucket = cache.info.name.clone();
let disk_path = self.path().to_string_lossy().to_string();
let source = match scan_mode {
HealScanMode::Deep => rustfs_scanner_contracts::metrics::ScannerWorkSource::Bitrot,
HealScanMode::Normal | HealScanMode::Unknown => rustfs_scanner_contracts::metrics::ScannerWorkSource::Usage,
HealScanMode::Deep => rustfs_scanner_metrics::metrics::ScannerWorkSource::Bitrot,
HealScanMode::Normal | HealScanMode::Unknown => rustfs_scanner_metrics::metrics::ScannerWorkSource::Usage,
};
global_metrics().record_scan_bucket_drive_start(source, &bucket, &disk_path);
let mut failure_guard = BucketDriveFailureGuard::new(source, &bucket, &disk_path);
+41
View File
@@ -26,12 +26,32 @@ use crate::{
ScannerPutObjReader, UNKNOWN_TIER, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
init_local_disks_with_instance_ctx, new_disk, path2_bucket_object_with_base_path,
};
use rustfs_concurrency::{
AdmissionState, WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshot, WorkloadAdmissionSnapshotProvider,
WorkloadClass,
};
use rustfs_filemeta::FileInfo;
use serial_test::serial;
use std::sync::Arc;
use temp_env::with_var;
use time::OffsetDateTime;
use uuid::Uuid;
#[derive(Clone)]
struct FixedWorkloadProvider {
snapshot: WorkloadAdmissionRegistrySnapshot,
}
impl WorkloadAdmissionSnapshotProvider for FixedWorkloadProvider {
fn workload_admission_snapshot(&self) -> WorkloadAdmissionRegistrySnapshot {
self.snapshot.clone()
}
}
fn install_scanner_workload_provider(snapshot: WorkloadAdmissionRegistrySnapshot) {
crate::set_scanner_workload_admission_snapshot_provider(Arc::new(FixedWorkloadProvider { snapshot }));
}
fn bucket_info(name: &str) -> BucketInfo {
BucketInfo {
name: name.to_string(),
@@ -1079,6 +1099,7 @@ async fn bucket_cache_pending_heal_reaches_cycle_maintenance_state() {
#[serial]
fn scanner_concurrency_limit_preserves_available_when_unconfigured() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
assert_eq!(scanner_concurrency_limit(0, 4), 4);
}
@@ -1086,6 +1107,7 @@ fn scanner_concurrency_limit_preserves_available_when_unconfigured() {
#[serial]
fn scanner_concurrency_limit_caps_to_configured_value() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
assert_eq!(scanner_concurrency_limit(2, 4), 2);
}
@@ -1093,6 +1115,7 @@ fn scanner_concurrency_limit_caps_to_configured_value() {
#[serial]
fn scanner_concurrency_limit_never_exceeds_available_work() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
assert_eq!(scanner_concurrency_limit(8, 4), 4);
}
@@ -1100,6 +1123,7 @@ fn scanner_concurrency_limit_never_exceeds_available_work() {
#[serial]
fn scanner_concurrency_limit_handles_no_available_work() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
assert_eq!(scanner_concurrency_limit(2, 0), 0);
}
@@ -1107,16 +1131,33 @@ fn scanner_concurrency_limit_handles_no_available_work() {
#[serial]
fn scanner_concurrency_limit_yields_to_foreground_reads() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
crate::set_foreground_read_activity(8);
assert_eq!(scanner_concurrency_limit(0, 4), 1);
assert_eq!(scanner_concurrency_limit(3, 4), 1);
crate::reset_foreground_read_activity_for_test();
}
#[test]
#[serial]
fn scanner_concurrency_limit_yields_to_shared_foreground_pressure() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
install_scanner_workload_provider(WorkloadAdmissionRegistrySnapshot::new(vec![
WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundWrite, AdmissionState::Open).with_counts(Some(2), None, Some(16)),
]));
assert_eq!(scanner_concurrency_limit(0, 4), 1);
assert_eq!(scanner_concurrency_limit(3, 4), 1);
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
}
#[test]
#[serial]
fn scanner_concurrency_limit_yields_to_streaming_reads() {
crate::reset_foreground_read_activity_for_test();
crate::workload_admission::clear_scanner_workload_admission_snapshot_provider_for_test();
let _guard = crate::ForegroundReadGuard::new();
assert_eq!(scanner_concurrency_limit(0, 4), 1);
+18 -16
View File
@@ -20,7 +20,7 @@ use rustfs_config::{
DEFAULT_SCANNER_IDLE_MODE, DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS, ENV_SCANNER_IDLE_MODE, ENV_SCANNER_SPEED,
ENV_SCANNER_YIELD_EVERY_N_OBJECTS, ScannerSpeed,
};
use rustfs_scanner_contracts::metrics::global_metrics;
use rustfs_scanner_metrics::metrics::global_metrics;
use tokio::time::Duration;
const MIN_SLEEP: Duration = Duration::from_millis(1);
@@ -29,8 +29,8 @@ const SCANNER_SPEED_FAST: u8 = 1;
const SCANNER_SPEED_DEFAULT: u8 = 2;
const SCANNER_SPEED_SLOW: u8 = 3;
const SCANNER_SPEED_SLOWEST: u8 = 4;
const FOREGROUND_READ_BACKOFF_PER_REQUEST_MS: u64 = 10;
const FOREGROUND_READ_BACKOFF_MAX_MS: u64 = 250;
const FOREGROUND_WORKLOAD_BACKOFF_PER_REQUEST_MS: u64 = 10;
const FOREGROUND_WORKLOAD_BACKOFF_MAX_MS: u64 = 250;
static SCANNER_DEFAULT_SPEED_PRESET: AtomicU8 = AtomicU8::new(SCANNER_SPEED_DEFAULT);
@@ -78,15 +78,15 @@ pub(crate) fn scanner_yield_every_n_objects() -> u64 {
rustfs_utils::get_env_u64(ENV_SCANNER_YIELD_EVERY_N_OBJECTS, DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS)
}
fn foreground_read_backoff_duration(active_reads: u64) -> Duration {
if active_reads == 0 {
fn foreground_workload_backoff_duration(active_foreground_workloads: u64) -> Duration {
if active_foreground_workloads == 0 {
return Duration::ZERO;
}
Duration::from_millis(
active_reads
.saturating_mul(FOREGROUND_READ_BACKOFF_PER_REQUEST_MS)
.min(FOREGROUND_READ_BACKOFF_MAX_MS),
active_foreground_workloads
.saturating_mul(FOREGROUND_WORKLOAD_BACKOFF_PER_REQUEST_MS)
.min(FOREGROUND_WORKLOAD_BACKOFF_MAX_MS),
)
}
@@ -147,14 +147,15 @@ impl DynamicSleeper {
}
let (factor, max_sleep) = self.read_params();
if factor == 0.0 || max_sleep.is_zero() {
let foreground_sleep = foreground_read_backoff_duration(crate::current_foreground_read_activity());
let foreground_sleep =
foreground_workload_backoff_duration(crate::workload_admission::foreground_workload_activity());
if !foreground_sleep.is_zero() {
tokio::time::sleep(foreground_sleep).await;
global_metrics().record_scanner_throttle_sleep(foreground_sleep);
}
return;
}
let foreground_sleep = foreground_read_backoff_duration(crate::current_foreground_read_activity());
let foreground_sleep = foreground_workload_backoff_duration(crate::workload_admission::foreground_workload_activity());
let sleep_dur = Duration::from_secs_f64(MIN_SLEEP.as_secs_f64() * factor)
.min(max_sleep)
.max(foreground_sleep);
@@ -235,7 +236,8 @@ impl SleepTimer {
}
let (factor, max_sleep) = self.sleeper.read_params();
if factor == 0.0 || max_sleep.is_zero() {
let foreground_sleep = foreground_read_backoff_duration(crate::current_foreground_read_activity());
let foreground_sleep =
foreground_workload_backoff_duration(crate::workload_admission::foreground_workload_activity());
if !foreground_sleep.is_zero() {
tokio::time::sleep(foreground_sleep).await;
global_metrics().record_scanner_throttle_sleep(foreground_sleep);
@@ -243,7 +245,7 @@ impl SleepTimer {
return;
}
let elapsed = self.start.elapsed();
let foreground_sleep = foreground_read_backoff_duration(crate::current_foreground_read_activity());
let foreground_sleep = foreground_workload_backoff_duration(crate::workload_admission::foreground_workload_activity());
let sleep_dur = Duration::from_secs_f64(elapsed.as_secs_f64() * factor)
.max(MIN_SLEEP)
.min(max_sleep)
@@ -304,10 +306,10 @@ mod tests {
}
#[test]
fn foreground_read_backoff_is_capped() {
assert_eq!(foreground_read_backoff_duration(0), Duration::ZERO);
assert_eq!(foreground_read_backoff_duration(1), Duration::from_millis(10));
assert_eq!(foreground_read_backoff_duration(80), Duration::from_millis(250));
fn foreground_workload_backoff_is_capped() {
assert_eq!(foreground_workload_backoff_duration(0), Duration::ZERO);
assert_eq!(foreground_workload_backoff_duration(1), Duration::from_millis(10));
assert_eq!(foreground_workload_backoff_duration(80), Duration::from_millis(250));
}
#[test]
+165
View File
@@ -0,0 +1,165 @@
// 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 std::sync::{Arc, LazyLock, RwLock};
use rustfs_concurrency::{AdmissionState, WorkloadAdmissionRegistrySnapshot, WorkloadAdmissionSnapshotProvider, WorkloadClass};
type WorkloadSnapshotProviderRef = Arc<dyn WorkloadAdmissionSnapshotProvider + Send + Sync>;
static SCANNER_WORKLOAD_ADMISSION_PROVIDER: LazyLock<RwLock<Option<WorkloadSnapshotProviderRef>>> =
LazyLock::new(|| RwLock::new(None));
pub fn set_scanner_workload_admission_snapshot_provider(provider: WorkloadSnapshotProviderRef) {
*SCANNER_WORKLOAD_ADMISSION_PROVIDER
.write()
.unwrap_or_else(|err| err.into_inner()) = Some(provider);
}
fn scanner_workload_admission_snapshot_provider() -> Option<WorkloadSnapshotProviderRef> {
SCANNER_WORKLOAD_ADMISSION_PROVIDER
.read()
.unwrap_or_else(|err| err.into_inner())
.clone()
}
#[cfg(test)]
pub(crate) fn clear_scanner_workload_admission_snapshot_provider_for_test() {
*SCANNER_WORKLOAD_ADMISSION_PROVIDER
.write()
.unwrap_or_else(|err| err.into_inner()) = None;
}
pub(crate) fn foreground_workload_activity() -> u64 {
let local_activity = crate::current_foreground_read_activity();
let Some(provider) = scanner_workload_admission_snapshot_provider() else {
return local_activity;
};
local_activity.max(foreground_activity_from_snapshot(&provider.workload_admission_snapshot()))
}
fn foreground_activity_from_snapshot(snapshot: &WorkloadAdmissionRegistrySnapshot) -> u64 {
[WorkloadClass::ForegroundRead, WorkloadClass::ForegroundWrite]
.into_iter()
.filter_map(|class| snapshot.get(class))
.map(|entry| {
entry
.active
.or_else(|| {
matches!(entry.state, AdmissionState::Saturated).then(|| entry.limit.filter(|limit| *limit > 0).unwrap_or(1))
})
.map(usize_to_u64_saturated)
.unwrap_or(0)
})
.max()
.unwrap_or(0)
}
fn usize_to_u64_saturated(value: usize) -> u64 {
u64::try_from(value).unwrap_or(u64::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
use rustfs_concurrency::{AdmissionState, WorkloadAdmissionSnapshot};
use serial_test::serial;
#[derive(Clone)]
struct FixedWorkloadProvider {
snapshot: WorkloadAdmissionRegistrySnapshot,
}
impl WorkloadAdmissionSnapshotProvider for FixedWorkloadProvider {
fn workload_admission_snapshot(&self) -> WorkloadAdmissionRegistrySnapshot {
self.snapshot.clone()
}
}
fn install_provider(snapshot: WorkloadAdmissionRegistrySnapshot) {
set_scanner_workload_admission_snapshot_provider(Arc::new(FixedWorkloadProvider { snapshot }));
}
#[test]
#[serial]
fn foreground_workload_activity_falls_back_to_local_read_activity() {
clear_scanner_workload_admission_snapshot_provider_for_test();
crate::reset_foreground_read_activity_for_test();
crate::set_foreground_read_activity(3);
assert_eq!(foreground_workload_activity(), 3);
crate::reset_foreground_read_activity_for_test();
}
#[test]
#[serial]
fn foreground_workload_activity_uses_shared_provider_counts() {
clear_scanner_workload_admission_snapshot_provider_for_test();
crate::reset_foreground_read_activity_for_test();
install_provider(WorkloadAdmissionRegistrySnapshot::new(vec![
WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Open).with_counts(
Some(5),
None,
Some(8),
),
WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundWrite, AdmissionState::Open).with_counts(
Some(2),
None,
Some(4),
),
]));
assert_eq!(foreground_workload_activity(), 5);
clear_scanner_workload_admission_snapshot_provider_for_test();
}
#[test]
#[serial]
fn foreground_workload_activity_treats_saturation_without_counts_as_pressure() {
clear_scanner_workload_admission_snapshot_provider_for_test();
crate::reset_foreground_read_activity_for_test();
install_provider(WorkloadAdmissionRegistrySnapshot::new(vec![
WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundRead, AdmissionState::Saturated).with_counts(
None,
None,
Some(7),
),
]));
assert_eq!(foreground_workload_activity(), 7);
clear_scanner_workload_admission_snapshot_provider_for_test();
}
#[test]
#[serial]
fn foreground_workload_activity_treats_zero_limit_saturation_as_pressure() {
clear_scanner_workload_admission_snapshot_provider_for_test();
crate::reset_foreground_read_activity_for_test();
install_provider(WorkloadAdmissionRegistrySnapshot::new(vec![
WorkloadAdmissionSnapshot::new(WorkloadClass::ForegroundWrite, AdmissionState::Saturated).with_counts(
None,
None,
Some(0),
),
]));
assert_eq!(foreground_workload_activity(), 1);
clear_scanner_workload_admission_snapshot_provider_for_test();
}
}