fix(scanner): harden lifecycle and tiering backlog (#3469)

* fix(scanner): expose lifecycle transition backlog

* fix(scanner): expose lifecycle expiry backlog

* fix(scanner): preserve lifecycle backlog states

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-06-15 15:01:21 +08:00
committed by GitHub
parent 57260c8314
commit dd6b4c35ad
7 changed files with 578 additions and 49 deletions
@@ -45,7 +45,9 @@ use futures::Future;
use http::HeaderMap;
use lazy_static::lazy_static;
use rustfs_common::heal_channel::rep_has_active_rules;
use rustfs_common::metrics::{IlmAction, Metrics, ScannerLifecycleTransitionStateUpdate, global_metrics};
use rustfs_common::metrics::{
IlmAction, Metrics, ScannerLifecycleExpiryStateUpdate, ScannerLifecycleTransitionStateUpdate, global_metrics,
};
use rustfs_config::{
DEFAULT_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_QUEUE_SEND_TIMEOUT_MS, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX,
DEFAULT_TRANSITION_WORKERS_CAP, ENV_TRANSITION_QUEUE_CAPACITY, ENV_TRANSITION_QUEUE_SEND_TIMEOUT_MS, ENV_TRANSITION_WORKERS,
@@ -110,6 +112,7 @@ const ENV_STALE_UPLOADS_CLEANUP_INTERVAL: &str = "RUSTFS_API_STALE_UPLOADS_CLEAN
const DEFAULT_STALE_UPLOADS_EXPIRY: StdDuration = StdDuration::from_secs(24 * 60 * 60);
const DEFAULT_STALE_UPLOADS_CLEANUP_INTERVAL: StdDuration = StdDuration::from_secs(6 * 60 * 60);
const DATE_EXPIRY_EXISTING_OBJECTS_GRACE_SECS: i64 = 5;
const EXPIRY_WORKER_QUEUE_CAPACITY: usize = 1000;
lazy_static! {
pub static ref GLOBAL_ExpiryState: Arc<RwLock<ExpiryState>> = ExpiryState::new();
@@ -157,7 +160,7 @@ fn is_immediate_transition_source(src: &LcEventSrc) -> bool {
fn record_scanner_lifecycle_enqueue_result(src: &LcEventSrc, count: u64, queued: bool) {
if matches!(src, LcEventSrc::Scanner) {
global_metrics().record_scanner_ilm_enqueue_result(count, queued);
global_metrics().record_scanner_expiry_enqueue_result(count, queued);
}
}
@@ -261,6 +264,8 @@ struct ExpiryStats {
missed_expiry_tasks: AtomicI64,
missed_freevers_tasks: AtomicI64,
missed_tier_journal_tasks: AtomicI64,
pending_tasks: AtomicI64,
active_tasks: AtomicI64,
workers: AtomicI64,
}
@@ -278,9 +283,91 @@ impl ExpiryStats {
self.missed_tier_journal_tasks.load(Ordering::SeqCst)
}
pub fn pending_tasks(&self) -> i64 {
self.pending_tasks.load(Ordering::SeqCst)
}
pub fn active_tasks(&self) -> i64 {
self.active_tasks.load(Ordering::SeqCst)
}
fn num_workers(&self) -> i64 {
self.workers.load(Ordering::SeqCst)
}
fn add_nonnegative(counter: &AtomicI64, delta: i64) {
let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| Some(current.saturating_add(delta).max(0)));
}
fn increment_missed_expiry_tasks(&self) {
Self::add_nonnegative(&self.missed_expiry_tasks, 1);
}
fn increment_missed_freevers_tasks(&self) {
Self::add_nonnegative(&self.missed_freevers_tasks, 1);
}
fn increment_missed_tier_journal_tasks(&self) {
Self::add_nonnegative(&self.missed_tier_journal_tasks, 1);
}
fn increment_pending_tasks(&self) {
Self::add_nonnegative(&self.pending_tasks, 1);
}
fn decrement_pending_tasks(&self) {
Self::add_nonnegative(&self.pending_tasks, -1);
}
fn increment_active_tasks(&self) {
Self::add_nonnegative(&self.active_tasks, 1);
}
fn decrement_active_tasks(&self) {
Self::add_nonnegative(&self.active_tasks, -1);
}
fn increment_workers(&self) {
Self::add_nonnegative(&self.workers, 1);
}
fn decrement_workers(&self) {
Self::add_nonnegative(&self.workers, -1);
}
fn scanner_expiry_state_update(&self) -> ScannerLifecycleExpiryStateUpdate {
let workers = nonnegative_i64_to_u64(self.num_workers());
ScannerLifecycleExpiryStateUpdate {
queue_capacity: workers.saturating_mul(usize_to_u64_saturated(EXPIRY_WORKER_QUEUE_CAPACITY)),
queued: nonnegative_i64_to_u64(self.pending_tasks()),
active: nonnegative_i64_to_u64(self.active_tasks()),
workers,
queue_missed: nonnegative_i64_to_u64(self.missed_tasks()),
}
}
fn record_scanner_expiry_state(&self) {
global_metrics().record_scanner_lifecycle_expiry_state(self.scanner_expiry_state_update());
}
}
struct ExpiryActiveTask {
stats: Arc<ExpiryStats>,
}
impl ExpiryActiveTask {
fn begin(stats: Arc<ExpiryStats>) -> Self {
stats.increment_active_tasks();
stats.record_scanner_expiry_state();
Self { stats }
}
}
impl Drop for ExpiryActiveTask {
fn drop(&mut self) {
self.stats.decrement_active_tasks();
self.stats.record_scanner_expiry_state();
}
}
pub trait ExpiryOp: 'static {
@@ -334,7 +421,7 @@ impl ExpiryOp for NewerNoncurrentTask {
pub struct ExpiryState {
tasks_tx: Vec<Sender<Option<ExpiryOpType>>>,
tasks_rx: Vec<Arc<tokio::sync::Mutex<Receiver<Option<ExpiryOpType>>>>>,
stats: Option<ExpiryStats>,
stats: Arc<ExpiryStats>,
}
impl ExpiryState {
@@ -343,41 +430,43 @@ impl ExpiryState {
Arc::new(RwLock::new(Self {
tasks_tx: vec![],
tasks_rx: vec![],
stats: Some(ExpiryStats {
stats: Arc::new(ExpiryStats {
missed_expiry_tasks: AtomicI64::new(0),
missed_freevers_tasks: AtomicI64::new(0),
missed_tier_journal_tasks: AtomicI64::new(0),
pending_tasks: AtomicI64::new(0),
active_tasks: AtomicI64::new(0),
workers: AtomicI64::new(0),
}),
}))
}
pub async fn pending_tasks(&self) -> usize {
let rxs = &self.tasks_rx;
if rxs.is_empty() {
return 0;
pub fn pending_tasks(&self) -> usize {
usize::try_from(self.stats.pending_tasks().max(0)).unwrap_or(usize::MAX)
}
async fn send_expiry_task(&self, wrkr: Sender<Option<ExpiryOpType>>, task: ExpiryOpType) -> bool {
self.stats.increment_pending_tasks();
let queued = wrkr.send(Some(task)).await.is_ok();
if !queued {
self.stats.decrement_pending_tasks();
}
let mut tasks = 0;
for rx in rxs.iter() {
tasks += rx.lock().await.len();
}
tasks
queued
}
pub async fn enqueue_tier_journal_entry(&mut self, je: &Jentry) -> Result<(), std::io::Error> {
let wrkr = self.get_worker_ch(je.op_hash());
if wrkr.is_none() {
*self.stats.as_mut().expect("stats lock").missed_tier_journal_tasks.get_mut() += 1;
self.stats.increment_missed_tier_journal_tasks();
self.stats.record_scanner_expiry_state();
return Ok(());
}
let wrkr = wrkr.expect("worker channel should exist after None check");
select! {
//_ -> GlobalContext.Done() => ()
_ = wrkr.send(Some(Box::new(je.clone()))) => (),
else => {
*self.stats.as_mut().expect("stats lock").missed_tier_journal_tasks.get_mut() += 1;
}
let queued = self.send_expiry_task(wrkr, Box::new(je.clone())).await;
if !queued {
self.stats.increment_missed_tier_journal_tasks();
}
self.stats.record_scanner_expiry_state();
Ok(())
}
@@ -385,17 +474,16 @@ impl ExpiryState {
let task = FreeVersionTask(oi);
let wrkr = self.get_worker_ch(task.op_hash());
if wrkr.is_none() {
*self.stats.as_mut().expect("stats lock").missed_freevers_tasks.get_mut() += 1;
self.stats.increment_missed_freevers_tasks();
self.stats.record_scanner_expiry_state();
return;
}
let wrkr = wrkr.expect("worker channel should exist after None check");
select! {
//_ -> GlobalContext.Done() => {}
_ = wrkr.send(Some(Box::new(task))) => (),
else => {
*self.stats.as_mut().expect("stats lock").missed_freevers_tasks.get_mut() += 1;
}
let queued = self.send_expiry_task(wrkr, Box::new(task)).await;
if !queued {
self.stats.increment_missed_freevers_tasks();
}
self.stats.record_scanner_expiry_state();
}
pub async fn enqueue_by_days(&mut self, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) -> bool {
@@ -406,16 +494,18 @@ impl ExpiryState {
};
let wrkr = self.get_worker_ch(task.op_hash());
if wrkr.is_none() {
*self.stats.as_mut().expect("stats lock").missed_expiry_tasks.get_mut() += 1;
self.stats.increment_missed_expiry_tasks();
record_scanner_lifecycle_enqueue_result(src, 1, false);
self.stats.record_scanner_expiry_state();
return false;
}
let wrkr = wrkr.expect("worker channel should exist after None check");
let queued = wrkr.send(Some(Box::new(task))).await.is_ok();
let queued = self.send_expiry_task(wrkr, Box::new(task)).await;
if !queued {
*self.stats.as_mut().expect("stats lock").missed_expiry_tasks.get_mut() += 1;
self.stats.increment_missed_expiry_tasks();
}
record_scanner_lifecycle_enqueue_result(src, 1, queued);
self.stats.record_scanner_expiry_state();
queued
}
@@ -438,16 +528,18 @@ impl ExpiryState {
};
let wrkr = self.get_worker_ch(task.op_hash());
if wrkr.is_none() {
*self.stats.as_mut().expect("stats lock").missed_expiry_tasks.get_mut() += 1;
self.stats.increment_missed_expiry_tasks();
record_scanner_lifecycle_enqueue_result(src, version_count, false);
self.stats.record_scanner_expiry_state();
return false;
}
let wrkr = wrkr.expect("worker channel should exist after None check");
let queued = wrkr.send(Some(Box::new(task))).await.is_ok();
let queued = self.send_expiry_task(wrkr, Box::new(task)).await;
if !queued {
*self.stats.as_mut().expect("stats lock").missed_expiry_tasks.get_mut() += 1;
self.stats.increment_missed_expiry_tasks();
}
record_scanner_lifecycle_enqueue_result(src, version_count, queued);
self.stats.record_scanner_expiry_state();
queued
}
@@ -459,7 +551,8 @@ impl ExpiryState {
}
pub fn increment_missed_tier_journal_tasks(&mut self) {
*self.stats.as_mut().expect("stats lock").missed_tier_journal_tasks.get_mut() += 1;
self.stats.increment_missed_tier_journal_tasks();
self.stats.record_scanner_expiry_state();
}
pub async fn resize_workers(n: usize, api: Arc<ECStore>) {
@@ -470,16 +563,17 @@ impl ExpiryState {
let mut state = GLOBAL_ExpiryState.write().await;
while state.tasks_tx.len() < n {
let (tx, rx) = mpsc::channel(1000);
let (tx, rx) = mpsc::channel(EXPIRY_WORKER_QUEUE_CAPACITY);
let api = api.clone();
let rx = Arc::new(tokio::sync::Mutex::new(rx));
let stats = Arc::clone(&state.stats);
state.tasks_tx.push(tx);
state.tasks_rx.push(rx.clone());
*state.stats.as_mut().expect("stats lock").workers.get_mut() += 1;
state.stats.increment_workers();
tokio::spawn(async move {
let mut rx = rx.lock().await;
//let mut expiry_state = GLOBAL_ExpiryState.read().await;
ExpiryState::worker(&mut rx, api).await;
ExpiryState::worker(&mut rx, api, stats).await;
});
}
@@ -489,12 +583,13 @@ impl ExpiryState {
worker.send(None).await.unwrap_or(());
state.tasks_tx.remove(l - 1);
state.tasks_rx.remove(l - 1);
*state.stats.as_mut().expect("stats lock").workers.get_mut() -= 1;
state.stats.decrement_workers();
l -= 1;
}
state.stats.record_scanner_expiry_state();
}
pub async fn worker(rx: &mut Receiver<Option<ExpiryOpType>>, api: Arc<ECStore>) {
async fn worker(rx: &mut Receiver<Option<ExpiryOpType>>, api: Arc<ECStore>, stats: Arc<ExpiryStats>) {
let cancel_token = crate::global::get_background_services_cancel_token().unwrap_or_else(|| {
static FALLBACK: std::sync::OnceLock<tokio_util::sync::CancellationToken> = std::sync::OnceLock::new();
FALLBACK.get_or_init(tokio_util::sync::CancellationToken::new)
@@ -525,6 +620,8 @@ impl ExpiryState {
return;
}
let v = v.expect("received None after None check");
stats.decrement_pending_tasks();
let _active_task = ExpiryActiveTask::begin(Arc::clone(&stats));
if v.as_any().is::<ExpiryTask>() {
let v = v.as_any().downcast_ref::<ExpiryTask>().expect("ExpiryTask downcast failed");
//debug!("lifecycle expiry worker received task: {:?}", v.obj_info);
@@ -803,6 +900,7 @@ impl TransitionState {
queue_full: nonnegative_i64_to_u64(Self::counter_value(&self.queue_full_tasks)),
queue_send_timeout: nonnegative_i64_to_u64(Self::counter_value(&self.queue_send_timeout_tasks)),
compensation_scheduled: nonnegative_i64_to_u64(Self::counter_value(&self.compensation_scheduled_tasks)),
compensation_pending: self.compensation_pending_tasks(),
compensation_running: nonnegative_i64_to_u64(Self::counter_value(&self.compensation_running_tasks)),
}
}
@@ -1002,6 +1100,13 @@ impl TransitionState {
Self::counter_value(&self.compensation_scheduled_tasks)
}
pub fn compensation_pending_tasks(&self) -> u64 {
match self.compensation_buckets.lock() {
Ok(scheduled) => usize_to_u64_saturated(scheduled.len()),
Err(poisoned) => usize_to_u64_saturated(poisoned.into_inner().len()),
}
}
pub fn compensation_running_tasks(&self) -> i64 {
Self::counter_value(&self.compensation_running_tasks)
}
@@ -2655,8 +2760,13 @@ mod tests {
let queued = state.enqueue_by_days(&object, &event, &LcEventSrc::Scanner).await;
assert!(!queued);
let stats = state.stats.as_ref().expect("expiry stats should exist");
assert_eq!(stats.missed_tasks(), 1);
assert_eq!(state.stats.missed_tasks(), 1);
let expiry = global_metrics().report().await.lifecycle_expiry;
assert_eq!(expiry.current_queue_capacity, 0);
assert_eq!(expiry.current_queued, 0);
assert_eq!(expiry.current_active, 0);
assert_eq!(expiry.current_workers, 0);
assert_eq!(expiry.queue_missed, 1);
}
#[tokio::test]
@@ -3060,6 +3170,17 @@ mod tests {
assert_eq!(state.compensation_scheduled_tasks(), 1);
}
#[tokio::test(flavor = "current_thread")]
async fn scanner_transition_state_reports_compensation_pending_buckets() {
let state = TransitionState::new_with_capacity(1);
assert_eq!(state.scanner_transition_state_update().compensation_pending, 0);
state.compensation_buckets.lock().unwrap().insert("bucket-a".to_string());
assert_eq!(state.scanner_transition_state_update().compensation_pending, 1);
state.compensation_buckets.lock().unwrap().insert("bucket-a".to_string());
assert_eq!(state.scanner_transition_state_update().compensation_pending, 1);
}
#[tokio::test]
#[serial]
async fn transition_state_init_honors_runtime_configured_worker_count() {
+29
View File
@@ -19,6 +19,7 @@ use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_madmin::metrics::{
DiskIOStats, DiskMetric, LastMinute as MadminLastMinute, NetDevLine, NetMetrics, RPCMetrics, RealtimeMetrics,
ScannerCheckpointReport as MadminScannerCheckpointReport,
ScannerLifecycleExpirySnapshot as MadminScannerLifecycleExpirySnapshot,
ScannerLifecycleTransitionSnapshot as MadminScannerLifecycleTransitionSnapshot,
ScannerMaintenanceControlSnapshot as MadminScannerMaintenanceControlSnapshot,
ScannerMaintenanceSourceSnapshot as MadminScannerMaintenanceSourceSnapshot, ScannerMetrics as MadminScannerMetrics,
@@ -173,12 +174,22 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
queue_full: metrics.lifecycle_transition.queue_full,
queue_send_timeout: metrics.lifecycle_transition.queue_send_timeout,
compensation_scheduled: metrics.lifecycle_transition.compensation_scheduled,
compensation_pending: metrics.lifecycle_transition.compensation_pending,
compensation_running: metrics.lifecycle_transition.compensation_running,
scanner_queued: metrics.lifecycle_transition.scanner_queued,
scanner_missed: metrics.lifecycle_transition.scanner_missed,
completed: metrics.lifecycle_transition.completed,
failed: metrics.lifecycle_transition.failed,
},
lifecycle_expiry: MadminScannerLifecycleExpirySnapshot {
current_queue_capacity: metrics.lifecycle_expiry.current_queue_capacity,
current_queued: metrics.lifecycle_expiry.current_queued,
current_active: metrics.lifecycle_expiry.current_active,
current_workers: metrics.lifecycle_expiry.current_workers,
queue_missed: metrics.lifecycle_expiry.queue_missed,
scanner_queued: metrics.lifecycle_expiry.scanner_queued,
scanner_missed: metrics.lifecycle_expiry.scanner_missed,
},
maintenance_control: MadminScannerMaintenanceControlSnapshot {
primary_control: metrics.maintenance_control.primary_control,
sources: metrics
@@ -566,6 +577,15 @@ mod test {
current_cycle_lifecycle_transition_actions: 3,
last_cycle_lifecycle_expiry_actions: 5,
last_cycle_lifecycle_transition_actions: 7,
lifecycle_expiry: rustfs_common::metrics::ScannerLifecycleExpirySnapshot {
current_queue_capacity: 16,
current_queued: 5,
current_active: 2,
current_workers: 4,
queue_missed: 3,
scanner_queued: 6,
scanner_missed: 2,
},
lifecycle_transition: rustfs_common::metrics::ScannerLifecycleTransitionSnapshot {
current_queue_capacity: 16,
current_queued: 5,
@@ -574,6 +594,7 @@ mod test {
queue_full: 3,
queue_send_timeout: 1,
compensation_scheduled: 2,
compensation_pending: 3,
compensation_running: 1,
scanner_queued: 6,
scanner_missed: 2,
@@ -583,6 +604,13 @@ mod test {
..Default::default()
});
assert_eq!(scanner.lifecycle_expiry.current_queue_capacity, 16);
assert_eq!(scanner.lifecycle_expiry.current_queued, 5);
assert_eq!(scanner.lifecycle_expiry.current_active, 2);
assert_eq!(scanner.lifecycle_expiry.current_workers, 4);
assert_eq!(scanner.lifecycle_expiry.queue_missed, 3);
assert_eq!(scanner.lifecycle_expiry.scanner_queued, 6);
assert_eq!(scanner.lifecycle_expiry.scanner_missed, 2);
assert_eq!(scanner.lifecycle_transition.current_queue_capacity, 16);
assert_eq!(scanner.lifecycle_transition.current_queued, 5);
assert_eq!(scanner.lifecycle_transition.current_active, 2);
@@ -590,6 +618,7 @@ mod test {
assert_eq!(scanner.lifecycle_transition.queue_full, 3);
assert_eq!(scanner.lifecycle_transition.queue_send_timeout, 1);
assert_eq!(scanner.lifecycle_transition.compensation_scheduled, 2);
assert_eq!(scanner.lifecycle_transition.compensation_pending, 3);
assert_eq!(scanner.lifecycle_transition.compensation_running, 1);
assert_eq!(scanner.lifecycle_transition.scanner_queued, 6);
assert_eq!(scanner.lifecycle_transition.scanner_missed, 2);