fix(ilm): keep expiry pending gauge balanced and drain overwrite tails (#7944)

This commit is contained in:
Chris
2026-09-17 08:38:02 +08:00
committed by GitHub
parent cd10981132
commit f018628374
2 changed files with 135 additions and 11 deletions
@@ -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<Option<ExpiryOpType>>, 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<RwLock<ExpiryStat
return false;
};
// Same ordering rule as `ExpiryState::send_expiry_task`: count first, so
// a worker that dequeues immediately cannot decrement before this
// increment lands.
stats.increment_pending_tasks();
let queued = wrkr.try_send(Some(Box::new(task))).is_ok();
if !queued {
stats.decrement_pending_tasks();
stats.increment_missed_freevers_tasks();
} else {
stats.increment_pending_tasks();
}
stats.record_scanner_expiry_state();
queued
@@ -7569,6 +7582,65 @@ mod tests {
assert!(recovery_notify.notified().now_or_never().is_some());
}
#[tokio::test]
async fn expiry_pending_gauge_survives_dequeue_landing_before_enqueue_accounting() {
// A worker may dequeue and decrement before the enqueuing side records
// its increment. The gauge must return to zero afterwards instead of
// clamping the transient -1 away and reporting a phantom pending task
// that no idle check can ever drain (rustfs#7921).
let state = ExpiryState::new();
let stats = Arc::clone(&state.read().await.stats);
stats.decrement_pending_tasks();
stats.increment_pending_tasks();
assert_eq!(stats.pending_tasks(), 0);
assert_eq!(state.read().await.pending_tasks(), 0);
stats.decrement_active_tasks();
stats.increment_active_tasks();
assert_eq!(stats.active_tasks(), 0);
assert_eq!(state.read().await.active_tasks(), 0);
}
#[tokio::test]
async fn free_version_enqueue_rolls_back_pending_when_queue_full() {
// Single-threaded, so this cannot observe the count-before-publish
// ordering itself; it pins the rollback that ordering requires: a
// rejected send must not leave its speculative increment behind.
let state = ExpiryState::new_with_unconsumed_worker_channel(1);
let oi = ObjectInfo {
bucket: "bucket".to_string(),
name: "object".to_string(),
transitioned_object: TransitionedObject {
name: "remote/object".to_string(),
version_id: "remote-version".to_string(),
tier: "WARM".to_string(),
free_version: true,
..Default::default()
},
..Default::default()
};
assert!(state.read().await.enqueue_free_version(oi.clone()));
assert_eq!(state.read().await.stats.pending_tasks(), 1);
// The single-slot queue is full for both enqueue paths.
assert!(!enqueue_recovered_free_version_with_state(&state, oi.clone()).await);
assert_eq!(state.read().await.stats.pending_tasks(), 1);
assert_eq!(state.read().await.stats.missed_free_vers_tasks(), 1);
assert!(!state.read().await.enqueue_free_version(oi));
assert_eq!(state.read().await.stats.pending_tasks(), 1);
assert_eq!(state.read().await.stats.missed_free_vers_tasks(), 2);
// Draining the one real task returns the gauge to zero.
let receiver = state.read().await.tasks_rx[0].clone();
let task = receiver.lock().await.recv().await.expect("queued task");
assert!(task.is_some());
state.read().await.stats.decrement_pending_tasks();
assert_eq!(state.read().await.pending_tasks(), 0);
}
#[tokio::test]
async fn enqueue_recovered_free_version_reports_false_without_worker_channel() {
let state = ExpiryState::new();
+55 -3
View File
@@ -2905,8 +2905,17 @@ mod tests {
if !version.is_empty() {
insert_str(&mut metadata, SUFFIX_TRANSITIONED_VERSION_ID, version.to_string());
}
// The commit path acknowledges at write quorum and drains
// the remaining rename fan-out in the background. This
// fixture "crashes" the store right after the overwrite,
// so it must wait for that tail: a disk left with the old
// live transitioned source makes exact cleanup fail closed
// (a minority live owner still references the remote
// tuple) until heal repairs it, which this fixture never
// runs (rustfs#7921).
let options = ObjectOptions {
version_suspended: suspended,
write_completion: crate::object_api::WriteCompletion::TailDrained,
..Default::default()
};
store
@@ -2978,6 +2987,7 @@ mod tests {
assert_eq!(free.len(), 1, "{state:?}, suspended={suspended}, copy={self_copy}");
assert_eq!(free[0].transitioned_objname, remote);
assert_eq!(free[0].transition_version_state, state);
assert_every_disk_holds_only_the_cleanup_owner(&set, &bucket, object, &remote).await;
assert!(backend.contains(&remote).await, "commit must not delete remote bytes before cleanup");
let removed_before = backend.remove_count().await;
@@ -12755,10 +12765,45 @@ mod tests {
}
}
/// Every physical copy must carry the replacement plus its tier
/// free-version owner, and none may still hold the pre-overwrite live
/// transitioned source. A stale minority copy is exactly what a lost
/// early-ACK rename tail leaves behind, and exact cleanup refuses to
/// delete remote bytes while such a live reference exists.
#[cfg(feature = "test-util")]
async fn assert_every_disk_holds_only_the_cleanup_owner(
set: &crate::set_disk::SetDisks,
bucket: &str,
object: &str,
remote: &str,
) {
use crate::disk::DiskAPI as _;
let disk_object = rustfs_utils::path::encode_dir_object(object);
for (index, disk) in set.disk_inventory().await.into_iter().enumerate() {
let disk = disk.unwrap_or_else(|| panic!("disk{index} should be online"));
let raw = disk
.read_xl(bucket, &disk_object, false)
.await
.unwrap_or_else(|err| panic!("disk{index} xl.meta should be readable after the overwrite: {err:?}"));
let versions = rustfs_filemeta::FileMeta::load(&raw.buf)
.and_then(|meta| meta.get_all_file_info_versions(bucket, object, true))
.unwrap_or_else(|err| panic!("disk{index} xl.meta should decode: {err:?}"));
let all: Vec<_> = 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")]