mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-26 05:56:50 +00:00
fix(scanner): clarify follower status
This commit is contained in:
@@ -54,7 +54,7 @@ use rustfs_config::{
|
|||||||
};
|
};
|
||||||
use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS};
|
use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS};
|
||||||
use rustfs_data_usage::observed_data_usage_is_newer;
|
use rustfs_data_usage::observed_data_usage_is_newer;
|
||||||
use rustfs_lock::NamespaceLockGuard;
|
use rustfs_lock::{NamespaceLockGuard, error::LockError};
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use sha2::{Digest as _, Sha256};
|
use sha2::{Digest as _, Sha256};
|
||||||
use tokio::sync::{Notify, mpsc};
|
use tokio::sync::{Notify, mpsc};
|
||||||
@@ -184,6 +184,8 @@ fn notify_scanner_cycle_state_persist_test_hook(leader_epoch: u64) {
|
|||||||
#[derive(Clone, Copy, Debug, Serialize)]
|
#[derive(Clone, Copy, Debug, Serialize)]
|
||||||
#[non_exhaustive]
|
#[non_exhaustive]
|
||||||
pub struct ScannerCycleScheduleStatus {
|
pub struct ScannerCycleScheduleStatus {
|
||||||
|
execution_role: &'static str,
|
||||||
|
effective_interval_available: bool,
|
||||||
effective_interval_seconds: u64,
|
effective_interval_seconds: u64,
|
||||||
clean_idle_backoff_enabled: bool,
|
clean_idle_backoff_enabled: bool,
|
||||||
clean_idle_backoff_multiplier: u64,
|
clean_idle_backoff_multiplier: u64,
|
||||||
@@ -194,6 +196,8 @@ pub struct ScannerCycleScheduleStatus {
|
|||||||
impl Default for ScannerCycleScheduleStatus {
|
impl Default for ScannerCycleScheduleStatus {
|
||||||
fn default() -> Self {
|
fn default() -> Self {
|
||||||
Self {
|
Self {
|
||||||
|
execution_role: "unknown",
|
||||||
|
effective_interval_available: false,
|
||||||
effective_interval_seconds: 0,
|
effective_interval_seconds: 0,
|
||||||
clean_idle_backoff_enabled: false,
|
clean_idle_backoff_enabled: false,
|
||||||
clean_idle_backoff_multiplier: 1,
|
clean_idle_backoff_multiplier: 1,
|
||||||
@@ -230,6 +234,8 @@ fn record_scanner_cycle_schedule(
|
|||||||
.write()
|
.write()
|
||||||
.unwrap_or_else(|poisoned| poisoned.into_inner());
|
.unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
*schedule = ScannerCycleScheduleStatus {
|
*schedule = ScannerCycleScheduleStatus {
|
||||||
|
execution_role: "leader",
|
||||||
|
effective_interval_available: true,
|
||||||
effective_interval_seconds,
|
effective_interval_seconds,
|
||||||
clean_idle_backoff_enabled,
|
clean_idle_backoff_enabled,
|
||||||
clean_idle_backoff_multiplier: clean_idle_backoff_multiplier.max(1),
|
clean_idle_backoff_multiplier: clean_idle_backoff_multiplier.max(1),
|
||||||
@@ -238,8 +244,30 @@ fn record_scanner_cycle_schedule(
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn record_scanner_cycle_schedule_role(execution_role: &'static str) {
|
||||||
|
let mut schedule = SCANNER_CYCLE_SCHEDULE
|
||||||
|
.write()
|
||||||
|
.unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||||
|
*schedule = ScannerCycleScheduleStatus {
|
||||||
|
execution_role,
|
||||||
|
..ScannerCycleScheduleStatus::default()
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
fn reset_scanner_cycle_schedule() {
|
fn reset_scanner_cycle_schedule() {
|
||||||
record_scanner_cycle_schedule(Duration::ZERO, false, 1, false, 0);
|
record_scanner_cycle_schedule_role("unknown");
|
||||||
|
}
|
||||||
|
|
||||||
|
enum ScannerLeaderLockFailure<'a> {
|
||||||
|
Contended,
|
||||||
|
Failed(&'a LockError),
|
||||||
|
}
|
||||||
|
|
||||||
|
fn classify_scanner_leader_lock_failure(error: &LockError) -> ScannerLeaderLockFailure<'_> {
|
||||||
|
match error {
|
||||||
|
LockError::Timeout { .. } => ScannerLeaderLockFailure::Contended,
|
||||||
|
error => ScannerLeaderLockFailure::Failed(error),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Returns the base cycle interval.
|
/// Returns the base cycle interval.
|
||||||
@@ -2041,6 +2069,7 @@ async fn run_data_scanner_with_maintenance_state(
|
|||||||
let mut guard = match storeapi.new_ns_lock(RUSTFS_META_BUCKET, "leader.lock").await {
|
let mut guard = match storeapi.new_ns_lock(RUSTFS_META_BUCKET, "leader.lock").await {
|
||||||
Ok(ns_lock) => match ns_lock.get_write_lock_quiet(get_lock_acquire_timeout()).await {
|
Ok(ns_lock) => match ns_lock.get_write_lock_quiet(get_lock_acquire_timeout()).await {
|
||||||
Ok(guard) => {
|
Ok(guard) => {
|
||||||
|
record_scanner_cycle_schedule_role("leader");
|
||||||
record_scanner_leader_lock_state("acquired");
|
record_scanner_leader_lock_state("acquired");
|
||||||
global_metrics().record_scanner_leader_liveness("acquired", true, "").await;
|
global_metrics().record_scanner_leader_liveness("acquired", true, "").await;
|
||||||
debug!(
|
debug!(
|
||||||
@@ -2055,20 +2084,38 @@ async fn run_data_scanner_with_maintenance_state(
|
|||||||
guard
|
guard
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
record_scanner_leader_lock_state("contended");
|
match classify_scanner_leader_lock_failure(&e) {
|
||||||
global_metrics()
|
ScannerLeaderLockFailure::Contended => {
|
||||||
.record_scanner_leader_liveness("contended", false, e.to_string())
|
record_scanner_cycle_schedule_role("follower");
|
||||||
.await;
|
record_scanner_leader_lock_state("contended");
|
||||||
debug!(
|
global_metrics().record_scanner_leader_liveness("contended", false, "").await;
|
||||||
target: "rustfs::scanner",
|
debug!(
|
||||||
event = EVENT_SCANNER_LOCK_STATE,
|
target: "rustfs::scanner",
|
||||||
component = LOG_COMPONENT_SCANNER,
|
event = EVENT_SCANNER_LOCK_STATE,
|
||||||
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
component = LOG_COMPONENT_SCANNER,
|
||||||
lock_name = "leader.lock",
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
||||||
state = "contended",
|
lock_name = "leader.lock",
|
||||||
error = ?e,
|
state = "contended",
|
||||||
"Scanner leader lock contended"
|
"Scanner leader lock contended"
|
||||||
);
|
);
|
||||||
|
}
|
||||||
|
ScannerLeaderLockFailure::Failed(error) => {
|
||||||
|
record_scanner_leader_lock_state("acquire_failed");
|
||||||
|
global_metrics()
|
||||||
|
.record_scanner_leader_liveness("acquire_failed", false, error.to_string())
|
||||||
|
.await;
|
||||||
|
error!(
|
||||||
|
target: "rustfs::scanner",
|
||||||
|
event = EVENT_SCANNER_LOCK_STATE,
|
||||||
|
component = LOG_COMPONENT_SCANNER,
|
||||||
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
||||||
|
lock_name = "leader.lock",
|
||||||
|
state = "acquire_failed",
|
||||||
|
error = %error,
|
||||||
|
"Scanner leader lock acquisition failed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -4626,14 +4626,24 @@ fn scanner_cycle_schedule_status_reports_effective_backoff() {
|
|||||||
|
|
||||||
let status = scanner_cycle_schedule_status();
|
let status = scanner_cycle_schedule_status();
|
||||||
|
|
||||||
|
assert_eq!(status.execution_role, "leader");
|
||||||
|
assert!(status.effective_interval_available);
|
||||||
assert_eq!(status.effective_interval_seconds, 86_401);
|
assert_eq!(status.effective_interval_seconds, 86_401);
|
||||||
assert!(status.clean_idle_backoff_enabled);
|
assert!(status.clean_idle_backoff_enabled);
|
||||||
assert_eq!(status.clean_idle_backoff_multiplier, 2_048);
|
assert_eq!(status.clean_idle_backoff_multiplier, 2_048);
|
||||||
assert!(status.superseded_retry_backoff_enabled);
|
assert!(status.superseded_retry_backoff_enabled);
|
||||||
assert_eq!(status.superseded_cycles, 7);
|
assert_eq!(status.superseded_cycles, 7);
|
||||||
|
|
||||||
|
record_scanner_cycle_schedule_role("follower");
|
||||||
|
let status = scanner_cycle_schedule_status();
|
||||||
|
assert_eq!(status.execution_role, "follower");
|
||||||
|
assert!(!status.effective_interval_available);
|
||||||
|
assert_eq!(status.effective_interval_seconds, 0);
|
||||||
|
|
||||||
reset_scanner_cycle_schedule();
|
reset_scanner_cycle_schedule();
|
||||||
let status = scanner_cycle_schedule_status();
|
let status = scanner_cycle_schedule_status();
|
||||||
|
assert_eq!(status.execution_role, "unknown");
|
||||||
|
assert!(!status.effective_interval_available);
|
||||||
assert_eq!(status.effective_interval_seconds, 0);
|
assert_eq!(status.effective_interval_seconds, 0);
|
||||||
assert!(!status.clean_idle_backoff_enabled);
|
assert!(!status.clean_idle_backoff_enabled);
|
||||||
assert_eq!(status.clean_idle_backoff_multiplier, 1);
|
assert_eq!(status.clean_idle_backoff_multiplier, 1);
|
||||||
@@ -4641,6 +4651,33 @@ fn scanner_cycle_schedule_status_reports_effective_backoff() {
|
|||||||
assert_eq!(status.superseded_cycles, 0);
|
assert_eq!(status.superseded_cycles, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn scanner_leader_lock_failure_classifies_only_timeout_as_expected_contention() {
|
||||||
|
let timeout = LockError::timeout(".rustfs.sys/leader.lock@latest", Duration::from_secs(5));
|
||||||
|
assert!(matches!(
|
||||||
|
classify_scanner_leader_lock_failure(&timeout),
|
||||||
|
ScannerLeaderLockFailure::Contended
|
||||||
|
));
|
||||||
|
|
||||||
|
let failures = [
|
||||||
|
LockError::internal("lock service unavailable"),
|
||||||
|
LockError::network(
|
||||||
|
"leader lock transport unavailable",
|
||||||
|
std::io::Error::new(std::io::ErrorKind::ConnectionRefused, "connection refused"),
|
||||||
|
),
|
||||||
|
LockError::QuorumNotReached {
|
||||||
|
required: 3,
|
||||||
|
achieved: 1,
|
||||||
|
},
|
||||||
|
];
|
||||||
|
for failure in &failures {
|
||||||
|
assert!(matches!(
|
||||||
|
classify_scanner_leader_lock_failure(failure),
|
||||||
|
ScannerLeaderLockFailure::Failed(_)
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn clean_idle_backoff_resets_for_non_idle_work() {
|
fn clean_idle_backoff_resets_for_non_idle_work() {
|
||||||
let base_interval = Duration::from_secs(60);
|
let base_interval = Duration::from_secs(60);
|
||||||
|
|||||||
@@ -389,6 +389,8 @@ mod tests {
|
|||||||
);
|
);
|
||||||
|
|
||||||
let encoded = serde_json::to_value(response).expect("scanner status should serialize");
|
let encoded = serde_json::to_value(response).expect("scanner status should serialize");
|
||||||
|
assert_eq!(encoded["cycle_schedule"]["execution_role"], "unknown");
|
||||||
|
assert_eq!(encoded["cycle_schedule"]["effective_interval_available"], false);
|
||||||
assert_eq!(encoded["cycle_schedule"]["effective_interval_seconds"], 0);
|
assert_eq!(encoded["cycle_schedule"]["effective_interval_seconds"], 0);
|
||||||
assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_enabled"], false);
|
assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_enabled"], false);
|
||||||
assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_multiplier"], 1);
|
assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_multiplier"], 1);
|
||||||
|
|||||||
Reference in New Issue
Block a user