feat(metrics): expose deferred usage freshness (#6449)

This commit is contained in:
cxymds
2026-08-23 19:31:01 +08:00
committed by GitHub
parent b2e60be647
commit a8e4b67d99
8 changed files with 331 additions and 25 deletions
+104
View File
@@ -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<String>,
scanner_usage_last_publication_reason: Mutex<String>,
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<String>,
scanner_source_work: Vec<ScannerSourceWorkCounters>,
current_scan_cycle_source_work_start: Vec<ScannerSourceWorkCounters>,
last_scan_cycle_source_work: Vec<ScannerSourceWorkCounters>,
@@ -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<String>) {
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<String>) {
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]
@@ -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]
+75
View File
@@ -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();
+5
View File
@@ -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);
+1
View File
@@ -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]
+15 -1
View File
@@ -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;
}
}