From a8e4b67d9972efc2cf11e145b85b8c5deb2f6185 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 23 Aug 2026 19:31:01 +0800 Subject: [PATCH] feat(metrics): expose deferred usage freshness (#6449) --- crates/common/src/metrics.rs | 104 ++++++++++++++++ .../ecstore/src/services/metrics_realtime.rs | 24 ++++ crates/madmin/src/metrics.rs | 75 ++++++++++++ crates/scanner/src/scanner.rs | 5 + crates/scanner/src/scanner/tests.rs | 1 + crates/scanner/src/scanner/usage_store.rs | 16 ++- rustfs/src/admin/handlers/cluster_snapshot.rs | 115 ++++++++++++++---- rustfs/src/cluster_snapshot.rs | 16 +++ 8 files changed, 331 insertions(+), 25 deletions(-) diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index 813e39ac3..32cb1b18a 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -918,7 +918,15 @@ pub struct Metrics { scanner_dirty_usage_last_cycle_dirty_buckets: AtomicU64, scanner_dirty_usage_last_cycle_cleared_buckets: AtomicU64, scanner_usage_last_save_unix_secs: AtomicU64, + scanner_usage_last_durable_success_unix_secs: AtomicU64, + scanner_usage_last_publication_unix_secs: AtomicU64, + scanner_usage_last_publication_state: Mutex, + scanner_usage_last_publication_reason: Mutex, scanner_usage_last_save_result: AtomicU8, + scanner_usage_deferred_pending: AtomicBool, + scanner_usage_deferred_total: AtomicU64, + scanner_usage_last_deferred_unix_secs: AtomicU64, + scanner_usage_last_deferred_reason: Mutex, scanner_source_work: Vec, current_scan_cycle_source_work_start: Vec, last_scan_cycle_source_work: Vec, @@ -1216,6 +1224,22 @@ pub struct ScannerUsageFreshnessSnapshot { pub last_usage_save_unix_secs: u64, pub last_usage_save_result: String, pub last_usage_save_result_code: u64, + #[serde(default)] + pub last_durable_success_unix_secs: u64, + #[serde(default)] + pub last_publication_unix_secs: u64, + #[serde(default)] + pub last_publication_state: String, + #[serde(default)] + pub last_publication_reason: String, + #[serde(default)] + pub deferred_pending: bool, + #[serde(default)] + pub deferred_total: u64, + #[serde(default)] + pub last_deferred_unix_secs: u64, + #[serde(default)] + pub last_deferred_reason: String, } #[derive(Clone, Debug, Default, Serialize, Deserialize)] @@ -1945,7 +1969,15 @@ impl Metrics { scanner_dirty_usage_last_cycle_dirty_buckets: AtomicU64::new(0), scanner_dirty_usage_last_cycle_cleared_buckets: AtomicU64::new(0), scanner_usage_last_save_unix_secs: AtomicU64::new(0), + scanner_usage_last_durable_success_unix_secs: AtomicU64::new(0), + scanner_usage_last_publication_unix_secs: AtomicU64::new(0), + scanner_usage_last_publication_state: Mutex::new(String::new()), + scanner_usage_last_publication_reason: Mutex::new(String::new()), scanner_usage_last_save_result: AtomicU8::new(ScannerUsageSaveResult::Unknown as u8), + scanner_usage_deferred_pending: AtomicBool::new(false), + scanner_usage_deferred_total: AtomicU64::new(0), + scanner_usage_last_deferred_unix_secs: AtomicU64::new(0), + scanner_usage_last_deferred_reason: Mutex::new(String::new()), scanner_source_work: ScannerWorkSource::all() .iter() .map(|_| ScannerSourceWorkCounters::default()) @@ -2270,6 +2302,44 @@ impl Metrics { .store(unix_now_secs(), Ordering::Relaxed); } + /// Record an intentional retryable usage publication deferral separately + /// from the last durable save result. + pub fn record_scanner_usage_deferred(&self, reason: impl Into) { + let reason = reason.into(); + self.record_scanner_usage_publication("deferred", reason.clone()); + self.scanner_usage_deferred_pending.store(true, Ordering::Release); + self.scanner_usage_deferred_total.fetch_add(1, Ordering::Relaxed); + self.scanner_usage_last_deferred_unix_secs + .store(unix_now_secs(), Ordering::Relaxed); + let mut last_reason = match self.scanner_usage_last_deferred_reason.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + }; + *last_reason = reason; + } + + pub fn record_scanner_usage_durable_success(&self) { + self.record_scanner_usage_publication("success", ""); + self.scanner_usage_last_durable_success_unix_secs + .store(unix_now_secs(), Ordering::Relaxed); + self.scanner_usage_deferred_pending.store(false, Ordering::Release); + } + + pub fn record_scanner_usage_publication(&self, state: &str, reason: impl Into) { + self.scanner_usage_last_publication_unix_secs + .store(unix_now_secs(), Ordering::Relaxed); + let mut publication_state = match self.scanner_usage_last_publication_state.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + }; + *publication_state = state.to_string(); + let mut publication_reason = match self.scanner_usage_last_publication_reason.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + }; + *publication_reason = reason.into(); + } + pub fn record_scanner_source_work(&self, source: ScannerWorkSource, work: ScannerSourceWorkUpdate) { if let Some(counters) = self.scanner_source_work.get(source.index()) { counters.add(work); @@ -3292,6 +3362,23 @@ impl Metrics { last_usage_save_unix_secs: self.scanner_usage_last_save_unix_secs.load(Ordering::Relaxed), last_usage_save_result: usage_save_result.as_str().to_string(), last_usage_save_result_code: usage_save_result as u8 as u64, + last_durable_success_unix_secs: self.scanner_usage_last_durable_success_unix_secs.load(Ordering::Relaxed), + last_publication_unix_secs: self.scanner_usage_last_publication_unix_secs.load(Ordering::Relaxed), + last_publication_state: match self.scanner_usage_last_publication_state.lock() { + Ok(state) => state.clone(), + Err(poisoned) => poisoned.into_inner().clone(), + }, + last_publication_reason: match self.scanner_usage_last_publication_reason.lock() { + Ok(reason) => reason.clone(), + Err(poisoned) => poisoned.into_inner().clone(), + }, + deferred_pending: self.scanner_usage_deferred_pending.load(Ordering::Acquire), + deferred_total: self.scanner_usage_deferred_total.load(Ordering::Relaxed), + last_deferred_unix_secs: self.scanner_usage_last_deferred_unix_secs.load(Ordering::Relaxed), + last_deferred_reason: match self.scanner_usage_last_deferred_reason.lock() { + Ok(reason) => reason.clone(), + Err(poisoned) => poisoned.into_inner().clone(), + }, }; m.throttle_idle_mode_enabled = self.scanner_throttle_idle_mode_enabled.load(Ordering::Relaxed); m.throttle_sleep_factor = self.scanner_throttle_sleep_factor_micros.load(Ordering::Relaxed) as f64 / 1_000_000.0; @@ -4663,6 +4750,7 @@ mod tests { metrics.record_scanner_dirty_usage_cycle_snapshot(1); metrics.record_scanner_dirty_usage_cycle_clear(1, 1); metrics.record_scanner_usage_save_result(ScannerUsageSaveResult::Success); + metrics.record_scanner_usage_deferred("data_movement"); let report = metrics.report().await; @@ -4674,6 +4762,22 @@ mod tests { assert!(report.usage_freshness.last_usage_save_unix_secs > 0); assert_eq!(report.usage_freshness.last_usage_save_result, "success"); assert_eq!(report.usage_freshness.last_usage_save_result_code, 1); + assert!(report.usage_freshness.deferred_pending); + assert_eq!(report.usage_freshness.deferred_total, 1); + assert!(report.usage_freshness.last_deferred_unix_secs > 0); + assert_eq!(report.usage_freshness.last_deferred_reason, "data_movement"); + + metrics.record_scanner_usage_durable_success(); + let report = metrics.report().await; + assert!(!report.usage_freshness.deferred_pending); + assert_eq!(report.usage_freshness.deferred_total, 1); + assert!(report.usage_freshness.last_durable_success_unix_secs > 0); + assert_eq!(report.usage_freshness.last_publication_state, "success"); + + metrics.record_scanner_usage_publication("no_update", "no_update"); + let report = metrics.report().await; + assert_eq!(report.usage_freshness.last_publication_state, "no_update"); + assert_eq!(report.usage_freshness.last_publication_reason, "no_update"); } #[tokio::test] diff --git a/crates/ecstore/src/services/metrics_realtime.rs b/crates/ecstore/src/services/metrics_realtime.rs index 6e2e5cd7d..290504fcb 100644 --- a/crates/ecstore/src/services/metrics_realtime.rs +++ b/crates/ecstore/src/services/metrics_realtime.rs @@ -220,6 +220,14 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo last_usage_save_unix_secs: metrics.usage_freshness.last_usage_save_unix_secs, last_usage_save_result: metrics.usage_freshness.last_usage_save_result, last_usage_save_result_code: metrics.usage_freshness.last_usage_save_result_code, + last_durable_success_unix_secs: metrics.usage_freshness.last_durable_success_unix_secs, + last_publication_unix_secs: metrics.usage_freshness.last_publication_unix_secs, + last_publication_state: metrics.usage_freshness.last_publication_state, + last_publication_reason: metrics.usage_freshness.last_publication_reason, + deferred_pending: metrics.usage_freshness.deferred_pending, + deferred_total: metrics.usage_freshness.deferred_total, + last_deferred_unix_secs: metrics.usage_freshness.last_deferred_unix_secs, + last_deferred_reason: metrics.usage_freshness.last_deferred_reason, }, maintenance_control: MadminScannerMaintenanceControlSnapshot { primary_control: metrics.maintenance_control.primary_control, @@ -816,6 +824,14 @@ mod test { last_usage_save_unix_secs: 12, last_usage_save_result: "success".to_string(), last_usage_save_result_code: 1, + last_durable_success_unix_secs: 13, + last_publication_unix_secs: 14, + last_publication_state: "published".to_string(), + last_publication_reason: "complete".to_string(), + deferred_pending: true, + deferred_total: 15, + last_deferred_unix_secs: 16, + last_deferred_reason: "data_movement".to_string(), }, ..Default::default() }); @@ -828,6 +844,14 @@ mod test { assert_eq!(scanner.usage_freshness.last_usage_save_unix_secs, 12); assert_eq!(scanner.usage_freshness.last_usage_save_result, "success"); assert_eq!(scanner.usage_freshness.last_usage_save_result_code, 1); + assert_eq!(scanner.usage_freshness.last_durable_success_unix_secs, 13); + assert_eq!(scanner.usage_freshness.last_publication_unix_secs, 14); + assert_eq!(scanner.usage_freshness.last_publication_state, "published"); + assert_eq!(scanner.usage_freshness.last_publication_reason, "complete"); + assert!(scanner.usage_freshness.deferred_pending); + assert_eq!(scanner.usage_freshness.deferred_total, 15); + assert_eq!(scanner.usage_freshness.last_deferred_unix_secs, 16); + assert_eq!(scanner.usage_freshness.last_deferred_reason, "data_movement"); } #[test] diff --git a/crates/madmin/src/metrics.rs b/crates/madmin/src/metrics.rs index 4e55563a5..0ecd457f6 100644 --- a/crates/madmin/src/metrics.rs +++ b/crates/madmin/src/metrics.rs @@ -302,10 +302,28 @@ pub struct ScannerUsageFreshnessSnapshot { pub last_usage_save_result: String, #[serde(rename = "last_usage_save_result_code", default)] pub last_usage_save_result_code: u64, + #[serde(rename = "last_durable_success_unix_secs", default)] + pub last_durable_success_unix_secs: u64, + #[serde(rename = "last_publication_unix_secs", default)] + pub last_publication_unix_secs: u64, + #[serde(rename = "last_publication_state", default)] + pub last_publication_state: String, + #[serde(rename = "last_publication_reason", default)] + pub last_publication_reason: String, + #[serde(rename = "deferred_pending", default)] + pub deferred_pending: bool, + #[serde(rename = "deferred_total", default)] + pub deferred_total: u64, + #[serde(rename = "last_deferred_unix_secs", default)] + pub last_deferred_unix_secs: u64, + #[serde(rename = "last_deferred_reason", default)] + pub last_deferred_reason: String, } impl ScannerUsageFreshnessSnapshot { fn merge(&mut self, other: &Self) { + let self_deferred_state_at = self.last_deferred_unix_secs.max(self.last_durable_success_unix_secs); + let other_deferred_state_at = other.last_deferred_unix_secs.max(other.last_durable_success_unix_secs); self.dirty_pending_buckets = self.dirty_pending_buckets.saturating_add(other.dirty_pending_buckets); self.last_dirty_mark_unix_secs = self.last_dirty_mark_unix_secs.max(other.last_dirty_mark_unix_secs); self.last_dirty_clear_unix_secs = self.last_dirty_clear_unix_secs.max(other.last_dirty_clear_unix_secs); @@ -318,6 +336,22 @@ impl ScannerUsageFreshnessSnapshot { self.last_usage_save_result = other.last_usage_save_result.clone(); self.last_usage_save_result_code = other.last_usage_save_result_code; } + self.last_durable_success_unix_secs = self.last_durable_success_unix_secs.max(other.last_durable_success_unix_secs); + if other.last_publication_unix_secs > self.last_publication_unix_secs { + self.last_publication_unix_secs = other.last_publication_unix_secs; + self.last_publication_state = other.last_publication_state.clone(); + self.last_publication_reason = other.last_publication_reason.clone(); + } + self.deferred_total = self.deferred_total.saturating_add(other.deferred_total); + if other_deferred_state_at > self_deferred_state_at { + self.deferred_pending = other.deferred_pending; + } else if other_deferred_state_at == self_deferred_state_at { + self.deferred_pending |= other.deferred_pending; + } + if other.last_deferred_unix_secs > self.last_deferred_unix_secs { + self.last_deferred_unix_secs = other.last_deferred_unix_secs; + self.last_deferred_reason = other.last_deferred_reason.clone(); + } } } @@ -1566,6 +1600,47 @@ mod tests { assert_eq!(explicit_false.current_cycle_active, Some(false)); } + #[test] + fn usage_freshness_deferred_fields_are_backward_compatible_and_merge() { + let legacy: ScannerUsageFreshnessSnapshot = serde_json::from_value(serde_json::json!({ + "dirty_pending_buckets": 2, + "last_usage_save_result": "success" + })) + .expect("legacy usage freshness should decode"); + assert!(!legacy.deferred_pending); + assert_eq!(legacy.deferred_total, 0); + + let mut merged = ScannerUsageFreshnessSnapshot::default(); + merged.merge(&ScannerUsageFreshnessSnapshot { + deferred_pending: true, + deferred_total: 2, + last_deferred_unix_secs: 20, + last_deferred_reason: "data_movement".to_string(), + ..Default::default() + }); + merged.merge(&ScannerUsageFreshnessSnapshot { + deferred_total: 1, + last_deferred_unix_secs: 10, + last_deferred_reason: "older".to_string(), + ..Default::default() + }); + assert!(merged.deferred_pending); + assert_eq!(merged.deferred_total, 3); + assert_eq!(merged.last_deferred_unix_secs, 20); + assert_eq!(merged.last_deferred_reason, "data_movement"); + + merged.merge(&ScannerUsageFreshnessSnapshot { + deferred_pending: false, + last_durable_success_unix_secs: 30, + last_publication_unix_secs: 30, + last_publication_state: "success".to_string(), + ..Default::default() + }); + assert!(!merged.deferred_pending); + assert_eq!(merged.last_durable_success_unix_secs, 30); + assert_eq!(merged.last_publication_state, "success"); + } + #[test] fn scanner_metrics_merge_prefers_an_active_first_cycle() { let collected_at = Timestamp::now(); diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index c44ea8ae0..362eb79a7 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -1375,6 +1375,7 @@ async fn run_data_scanner_cycle_with_budget( Ok(result) => final_data_usage_publication_defer_reason(storeapi.as_ref(), result.status).await, Err(_) => Some(ScannerCycleDeferReason::ActivityBaselineUnavailable), }; + let publication_deferred = publication_defer_reason.is_some(); let publication_epoch = scan_result.as_ref().ok().and_then(ScannerCycleResult::publication_epoch); let budget_elapsed = cycle_budget.budget_elapsed() && !ctx.is_cancelled(); let usage_persist_outcome = match publication_defer_reason { @@ -1539,6 +1540,9 @@ async fn run_data_scanner_cycle_with_budget( state = "deferred", "Scanner cycle deferred before data usage publication" ); + if publication_deferred { + global_metrics().record_scanner_usage_deferred(reason.as_str()); + } emit_scan_cycle_deferred(cycle_start.elapsed()); mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Deferred(reason); @@ -1708,6 +1712,7 @@ async fn run_data_scanner_cycle_with_budget( state = "deferred", "Scanner cycle deferred before usage scanning began" ); + global_metrics().record_scanner_usage_deferred(reason.as_str()); emit_scan_cycle_deferred(cycle_start.elapsed()); mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Deferred(reason); diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 9e1d7cef8..b59cdc2d8 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -2788,6 +2788,7 @@ async fn test_deferred_usage_save_keeps_last_real_save_metric() { assert_eq!(after.last_usage_save_result, before.last_usage_save_result); assert_eq!(after.last_usage_save_result_code, before.last_usage_save_result_code); assert_eq!(after.last_usage_save_unix_secs, before.last_usage_save_unix_secs); + assert_eq!(after.deferred_total, before.deferred_total.saturating_add(1)); } #[tokio::test] diff --git a/crates/scanner/src/scanner/usage_store.rs b/crates/scanner/src/scanner/usage_store.rs index 9b211c1f5..24e5fd543 100644 --- a/crates/scanner/src/scanner/usage_store.rs +++ b/crates/scanner/src/scanner/usage_store.rs @@ -186,6 +186,7 @@ where state = "publication_blocked_before_reconcile", "Scanner data usage publication deferred by the pool-state fence" ); + global_metrics().record_scanner_usage_deferred(ScannerCycleDeferReason::DataMovement.as_str()); outcome = DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement); break; } @@ -288,6 +289,7 @@ where state = "reject_incomplete_snapshot", "Scanner refused to persist an incomplete data usage snapshot" ); + global_metrics().record_scanner_usage_publication("failed", "incomplete_snapshot"); global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Failed); outcome = DataUsagePersistOutcome::Failed; continue; @@ -306,6 +308,7 @@ where error = %e, "Scanner data usage encode failed" ); + global_metrics().record_scanner_usage_publication("failed", "encode_failed"); global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::EncodeFailed); outcome = DataUsagePersistOutcome::Failed; continue; @@ -549,6 +552,7 @@ where replace_bucket_usage_memory_from_info(&data_usage_info).await; } global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Success); + global_metrics().record_scanner_usage_durable_success(); outcome = DataUsagePersistOutcome::AlreadyDurable; } DataUsagePersistOutcome::PriorCycleDurable => { @@ -568,9 +572,17 @@ where invalidate_data_usage_snapshot_cache().await; } global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Success); + global_metrics().record_scanner_usage_durable_success(); outcome = DataUsagePersistOutcome::PriorCycleDurable; } - DataUsagePersistOutcome::Failed | DataUsagePersistOutcome::NoUpdate => { + DataUsagePersistOutcome::NoUpdate => { + global_metrics().record_scanner_usage_publication("no_update", "no_update"); + global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Failed); + outcome = DataUsagePersistOutcome::NoUpdate; + continue; + } + DataUsagePersistOutcome::Failed => { + global_metrics().record_scanner_usage_publication("failed", "save_failed"); global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Failed); outcome = DataUsagePersistOutcome::Failed; continue; @@ -579,6 +591,7 @@ where // A deferred publication is an intentional retryable state, not a // failed save. Keep the last real save result so admin freshness // reporting does not turn a pool-recovery fence into a false error. + global_metrics().record_scanner_usage_deferred(reason.as_str()); outcome = DataUsagePersistOutcome::Deferred(reason); break 'updates; } @@ -600,6 +613,7 @@ where replace_bucket_usage_memory_from_info(&data_usage_info).await; } global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Success); + global_metrics().record_scanner_usage_durable_success(); outcome = DataUsagePersistOutcome::Saved; } } diff --git a/rustfs/src/admin/handlers/cluster_snapshot.rs b/rustfs/src/admin/handlers/cluster_snapshot.rs index 3694d01a4..1646f1228 100644 --- a/rustfs/src/admin/handlers/cluster_snapshot.rs +++ b/rustfs/src/admin/handlers/cluster_snapshot.rs @@ -37,6 +37,9 @@ use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; use serde::Serialize; +use std::time::{SystemTime, UNIX_EPOCH}; + +const USAGE_DEFERRED_STALE_THRESHOLD_SECS: u64 = 300; pub fn register_cluster_snapshot_route(r: &mut S3Router) -> std::io::Result<()> { r.insert( @@ -242,7 +245,16 @@ pub(crate) struct ClusterUsageFreshnessStatus { pub last_usage_save_unix_secs: u64, pub last_usage_save_result: String, pub last_success_unix_secs: Option, + pub last_durable_success_unix_secs: u64, + pub last_publication_unix_secs: u64, + pub last_publication_state: String, + pub last_publication_reason: String, pub last_error: Option, + pub deferred_pending: bool, + pub deferred_total: u64, + pub last_deferred_unix_secs: u64, + pub last_deferred_reason: String, + pub deferred_age_secs: Option, } fn component_status(source: &'static str, status: CapabilityStatus) -> ClusterComponentStatus { @@ -701,32 +713,71 @@ fn summarize_listing_metacache(snapshot: &ClusterReadOnlySnapshot) -> ClusterLis fn summarize_usage_freshness(snapshot: &ClusterReadOnlySnapshot) -> ClusterUsageFreshnessStatus { let freshness = &snapshot.usage_freshness; - let (condition, status) = match freshness.last_usage_save_result.as_str() { - "success" if freshness.dirty_pending_buckets == 0 => ( - "healthy", - CapabilityStatus::supported().with_reason("usage cache was saved successfully and has no pending dirty buckets"), - ), - "success" | "" if freshness.dirty_pending_buckets > 0 => ( + let deferred_age_secs = (freshness.deferred_pending && freshness.last_deferred_unix_secs > 0).then(|| { + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |duration| duration.as_secs()); + now.saturating_sub(freshness.last_deferred_unix_secs) + }); + let deferred_stale = deferred_age_secs.is_some_and(|age| age > USAGE_DEFERRED_STALE_THRESHOLD_SECS); + let (condition, status) = if freshness.last_publication_state == "no_update" { + ( + "no_update", + CapabilityStatus::unknown().with_reason("usage cache publication produced no update"), + ) + } else if freshness.deferred_pending && deferred_stale { + ( "stale", - CapabilityStatus::unknown() - .with_reason(format!("usage cache has {} pending dirty buckets", freshness.dirty_pending_buckets)), - ), - "skipped_stale" => ( - "stale", - CapabilityStatus::unknown().with_reason("last usage cache save was skipped because scanner data was stale"), - ), - "failed" => ("degraded", CapabilityStatus::unknown().with_reason("last usage cache save failed")), - "encode_failed" => ( - "degraded", - CapabilityStatus::unknown().with_reason("last usage cache save failed during encoding"), - ), - _ => ( - "unknown", - CapabilityStatus::unknown().with_reason("no usage cache save result has been reported"), - ), + CapabilityStatus::unknown().with_reason(format!( + "usage cache publication has been deferred for {} seconds (threshold: {} seconds)", + deferred_age_secs.unwrap_or_default(), + USAGE_DEFERRED_STALE_THRESHOLD_SECS + )), + ) + } else if freshness.deferred_pending { + ( + "deferred", + CapabilityStatus::unknown().with_reason(if freshness.last_deferred_reason.is_empty() { + "usage cache publication is temporarily deferred" + } else { + freshness.last_deferred_reason.as_str() + }), + ) + } else { + match freshness.last_usage_save_result.as_str() { + "success" if freshness.dirty_pending_buckets == 0 => ( + "healthy", + CapabilityStatus::supported().with_reason("usage cache was saved successfully and has no pending dirty buckets"), + ), + "success" | "" if freshness.dirty_pending_buckets > 0 => ( + "stale", + CapabilityStatus::unknown() + .with_reason(format!("usage cache has {} pending dirty buckets", freshness.dirty_pending_buckets)), + ), + "skipped_stale" => ( + "stale", + CapabilityStatus::unknown().with_reason("last usage cache save was skipped because scanner data was stale"), + ), + "failed" => ("degraded", CapabilityStatus::unknown().with_reason("last usage cache save failed")), + "encode_failed" => ( + "degraded", + CapabilityStatus::unknown().with_reason("last usage cache save failed during encoding"), + ), + _ => ( + "unknown", + CapabilityStatus::unknown().with_reason("no usage cache save result has been reported"), + ), + } + }; + // Mixed-version nodes do not publish the additive durable timestamp yet; + // retain the legacy successful-save timestamp until all peers are upgraded. + let last_success_unix_secs = if freshness.last_durable_success_unix_secs > 0 { + Some(freshness.last_durable_success_unix_secs) + } else if freshness.last_usage_save_result == "success" && freshness.last_usage_save_unix_secs > 0 { + Some(freshness.last_usage_save_unix_secs) + } else { + None }; - let last_success_unix_secs = (freshness.last_usage_save_result == "success" && freshness.last_usage_save_unix_secs > 0) - .then_some(freshness.last_usage_save_unix_secs); let last_error = match freshness.last_usage_save_result.as_str() { "failed" | "skipped_stale" | "encode_failed" => Some(freshness.last_usage_save_result.clone()), _ => None, @@ -744,7 +795,16 @@ fn summarize_usage_freshness(snapshot: &ClusterReadOnlySnapshot) -> ClusterUsage last_usage_save_unix_secs: freshness.last_usage_save_unix_secs, last_usage_save_result: freshness.last_usage_save_result.clone(), last_success_unix_secs, + last_durable_success_unix_secs: freshness.last_durable_success_unix_secs, + last_publication_unix_secs: freshness.last_publication_unix_secs, + last_publication_state: freshness.last_publication_state.clone(), + last_publication_reason: freshness.last_publication_reason.clone(), last_error, + deferred_pending: freshness.deferred_pending, + deferred_total: freshness.deferred_total, + last_deferred_unix_secs: freshness.last_deferred_unix_secs, + last_deferred_reason: freshness.last_deferred_reason.clone(), + deferred_age_secs, } } @@ -1282,6 +1342,7 @@ mod tests { usage_freshness: ClusterUsageFreshnessSnapshot { dirty_pending_buckets: 0, last_usage_save_unix_secs: 456, + last_durable_success_unix_secs: 450, last_usage_save_result: "success".to_string(), last_usage_save_result_code: 1, ..Default::default() @@ -1295,6 +1356,12 @@ mod tests { assert_eq!(component.condition, "healthy"); assert_eq!(component.last_usage_save_unix_secs, 456); assert_eq!(component.last_usage_save_result, "success"); + assert_eq!(component.last_success_unix_secs, Some(450)); + + let mut legacy_snapshot = snapshot.clone(); + legacy_snapshot.usage_freshness.last_durable_success_unix_secs = 0; + let legacy_component = super::summarize_usage_freshness(&legacy_snapshot); + assert_eq!(legacy_component.last_success_unix_secs, Some(456)); } #[test] diff --git a/rustfs/src/cluster_snapshot.rs b/rustfs/src/cluster_snapshot.rs index 8112afaf9..e7535eb5d 100644 --- a/rustfs/src/cluster_snapshot.rs +++ b/rustfs/src/cluster_snapshot.rs @@ -66,6 +66,14 @@ pub struct ClusterUsageFreshnessSnapshot { pub last_usage_save_unix_secs: u64, pub last_usage_save_result: String, pub last_usage_save_result_code: u64, + pub last_durable_success_unix_secs: u64, + pub last_publication_unix_secs: u64, + pub last_publication_state: String, + pub last_publication_reason: String, + pub deferred_pending: bool, + pub deferred_total: u64, + pub last_deferred_unix_secs: u64, + pub last_deferred_reason: String, } impl From<&ScannerMetricsReport> for ClusterUsageFreshnessSnapshot { @@ -79,6 +87,14 @@ impl From<&ScannerMetricsReport> for ClusterUsageFreshnessSnapshot { last_usage_save_unix_secs: report.usage_freshness.last_usage_save_unix_secs, last_usage_save_result: report.usage_freshness.last_usage_save_result.clone(), last_usage_save_result_code: report.usage_freshness.last_usage_save_result_code, + last_durable_success_unix_secs: report.usage_freshness.last_durable_success_unix_secs, + last_publication_unix_secs: report.usage_freshness.last_publication_unix_secs, + last_publication_state: report.usage_freshness.last_publication_state.clone(), + last_publication_reason: report.usage_freshness.last_publication_reason.clone(), + deferred_pending: report.usage_freshness.deferred_pending, + deferred_total: report.usage_freshness.deferred_total, + last_deferred_unix_secs: report.usage_freshness.last_deferred_unix_secs, + last_deferred_reason: report.usage_freshness.last_deferred_reason.clone(), } } }