From f018628374dadfb2accb6f5a586109728d626bc5 Mon Sep 17 00:00:00 2001 From: Chris Date: Thu, 17 Sep 2026 08:38:02 +0800 Subject: [PATCH] fix(ilm): keep expiry pending gauge balanced and drain overwrite tails (#7944) --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 88 +++++++++++++++++-- crates/ecstore/src/store/init.rs | 58 +++++++++++- 2 files changed, 135 insertions(+), 11 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 8c4ac1322..7a04abb36 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -515,20 +515,26 @@ impl ExpiryStats { Self::add_nonnegative(&self.missed_tier_journal_tasks, 1); } + // The pending and active gauges are balanced by design: every increment + // has exactly one matching decrement. They must not saturate at zero on + // update, because a worker can dequeue (and decrement) before the + // enqueuing side has recorded its increment. Clamping that transient -1 + // to 0 turns the later +1 into a phantom task that never drains + // (rustfs#7921). Readers clamp negative snapshots instead. fn increment_pending_tasks(&self) { - Self::add_nonnegative(&self.pending_tasks, 1); + self.pending_tasks.fetch_add(1, Ordering::AcqRel); } fn decrement_pending_tasks(&self) { - Self::add_nonnegative(&self.pending_tasks, -1); + self.pending_tasks.fetch_sub(1, Ordering::AcqRel); } fn increment_active_tasks(&self) { - Self::add_nonnegative(&self.active_tasks, 1); + self.active_tasks.fetch_add(1, Ordering::AcqRel); } fn decrement_active_tasks(&self) { - Self::add_nonnegative(&self.active_tasks, -1); + self.active_tasks.fetch_sub(1, Ordering::AcqRel); } fn increment_workers(&self) { @@ -1074,9 +1080,13 @@ impl ExpiryState { } fn send_expiry_task(&self, wrkr: Sender>, task: ExpiryOpType) -> bool { + // Account for the task before a worker can observe it. The worker + // decrements on dequeue, so incrementing after `try_send` would let a + // fast dequeue run the gauge through zero first. + self.stats.increment_pending_tasks(); let queued = wrkr.try_send(Some(task)).is_ok(); - if queued { - self.stats.increment_pending_tasks(); + if !queued { + self.stats.decrement_pending_tasks(); } queued } @@ -1444,11 +1454,14 @@ async fn enqueue_recovered_free_version_with_state(state: &Arc = versions.versions.iter().chain(versions.free_versions.iter()).collect(); + let owners = all.iter().filter(|fi| fi.tier_free_version()).count(); + let live_sources = all + .iter() + .filter(|fi| !fi.tier_free_version() && fi.transitioned_objname == remote) + .count(); + assert_eq!(owners, 1, "disk{index} must hold exactly one cleanup owner after the drained overwrite"); + assert_eq!(live_sources, 0, "disk{index} must not retain the pre-overwrite live transitioned source"); + } + } + #[cfg(feature = "test-util")] async fn wait_for_expiry_workers_idle(store: &crate::store::ECStore) { let expiry_state = store.ctx.expiry_state(); - tokio::time::timeout(Duration::from_secs(30), async { + let idle = tokio::time::timeout(Duration::from_secs(30), async { loop { let idle = { let state = expiry_state.read().await; @@ -12778,8 +12823,15 @@ mod tests { } } }) - .await - .expect("lifecycle expiry workers should become idle"); + .await; + if idle.is_err() { + let state = expiry_state.read().await; + panic!( + "lifecycle expiry workers should become idle: pending={} active={}", + state.pending_tasks(), + state.active_tasks() + ); + } } #[cfg(feature = "test-util")]