From 47c5a3ab350d7e14b3cb7e0e62f29e21b4bd314e Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 15 Jul 2026 15:23:57 +0800 Subject: [PATCH] fix(ecstore): deduplicate concurrent ILM transition enqueues (#4839) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A single PUT can enqueue the same object for transition twice — immediately (enqueue_transition_after_write) and again from the startup/lifecycle compensation backfill. transition_object's own namespace lock is commented out, so the two attempts are not serialized as same-object work: the winner transitions the object and removes its local data, and the loser then reads that already-removed data via get_object_fileinfo(read_data = true) and logs a spurious "get_object_with_fileinfo err ... No such file" / lifecycle_tier_operation_failed. For a large (multipart) object both attempts can also race the source read, corrupting the transferred copy. Guard the transition queue with an in-flight set keyed by (bucket, object, version). queue_transition_task claims the key before sending; a duplicate enqueue is reported handled without queueing a second task. The worker releases the claim once the transition finishes, so a later lifecycle pass can still re-transition the object, and a failed enqueue releases it immediately. This is the sole path into the transition channel, so every enqueue is covered. Surfaced while validating the >128 MiB transition fix end-to-end in Docker (rustfs/rustfs#4811); tracked as its own defect since it is unrelated to the checksum/multipart-client bugs. Tests: same-version enqueues dedupe (direct reserve and via queue_transition_task with spare capacity), a distinct object still hits the full-queue path, and a released claim can be re-acquired. Refs: rustfs/backlog#1268 Co-authored-by: heihutu --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 126 +++++++++++++++++- 1 file changed, 124 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 1958d1c27..2f28db497 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -836,6 +836,12 @@ pub struct TransitionState { compensation_scheduled_tasks: AtomicI64, compensation_running_tasks: AtomicI64, compensation_buckets: Arc>>, + // (bucket, object, version) currently queued or being transitioned. A single + // PUT can enqueue the same object twice (immediate transition + startup + // compensation backfill); without this guard the duplicates run concurrently, + // and the loser races the winner's source cleanup — reading data the winner + // already removed and logging a spurious NotFound failure (rustfs/backlog#1268). + in_flight_transitions: Arc>>, last_day_stats: Arc>>, } @@ -869,10 +875,39 @@ impl TransitionState { compensation_scheduled_tasks: AtomicI64::new(0), compensation_running_tasks: AtomicI64::new(0), compensation_buckets: Arc::new(Mutex::new(HashSet::new())), + in_flight_transitions: Arc::new(Mutex::new(HashSet::new())), last_day_stats: Arc::new(Mutex::new(HashMap::new())), }) } + fn transition_key(oi: &ObjectInfo) -> (String, String, String) { + ( + oi.bucket.clone(), + oi.name.clone(), + oi.version_id.map(|v| v.to_string()).unwrap_or_default(), + ) + } + + /// Try to claim a transition for this exact (bucket, object, version). Returns + /// `false` when one is already queued or running, so the caller can skip the + /// duplicate enqueue (rustfs/backlog#1268). The claim is released in the + /// worker once the transition finishes, or here if the enqueue fails. + fn reserve_transition(&self, oi: &ObjectInfo) -> bool { + let key = Self::transition_key(oi); + match self.in_flight_transitions.lock() { + Ok(mut set) => set.insert(key), + Err(poisoned) => poisoned.into_inner().insert(key), + } + } + + fn release_transition(&self, oi: &ObjectInfo) { + let key = Self::transition_key(oi); + match self.in_flight_transitions.lock() { + Ok(mut set) => set.remove(&key), + Err(poisoned) => poisoned.into_inner().remove(&key), + }; + } + fn reserve_bucket_compensation(&self, bucket: &str) -> bool { let inserted = match self.compensation_buckets.lock() { Ok(mut scheduled) => scheduled.insert(bucket.to_string()), @@ -1046,6 +1081,15 @@ impl TransitionState { return false; } + // Deduplicate concurrent enqueues of the same object version. The claim is + // released by the worker after the transition finishes (rustfs/backlog#1268). + // This is a no-op, not a new enqueue, so it must not touch the enqueue + // counters (the first enqueue already counted this object). + if !self.reserve_transition(oi) { + self.record_scanner_transition_state(); + return true; + } + let task = TransitionTask { obj_info: oi.clone(), src: src.clone(), @@ -1084,6 +1128,9 @@ impl TransitionState { self.handle_immediate_enqueue_failure(oi, src, ImmediateEnqueueFailure::QueueClosed { timeout_ms: None }); } } + if !queued { + self.release_transition(oi); + } record_scanner_transition_enqueue_result(src, 1, queued); self.record_scanner_transition_state(); return queued; @@ -1113,6 +1160,9 @@ impl TransitionState { ); } } + if !queued { + self.release_transition(oi); + } record_scanner_transition_enqueue_result(src, 1, queued); self.record_scanner_transition_state(); queued @@ -1239,6 +1289,9 @@ impl TransitionState { emit_transition_complete_event(obj_info_for_event); } + // Release the dedup claim so a later lifecycle pass can + // re-transition this object if needed (rustfs/backlog#1268). + transition_state.release_transition(&task.obj_info); TransitionState::add_counter(&transition_state.active_tasks, -1); transition_state.record_scanner_transition_state(); } @@ -3244,8 +3297,16 @@ mod tests { ..Default::default() }; + // A distinct object fills past the capacity-1 queue and is reported as a + // missed enqueue. (Re-enqueuing the same object would instead be deduped; + // see transition_reserve_dedupes_same_object_version.) + let other = ObjectInfo { + bucket: "bucket".to_string(), + name: "object-2".to_string(), + ..Default::default() + }; let first = state.queue_transition_task(&object, &event, &LcEventSrc::Scanner).await; - let second = state.queue_transition_task(&object, &event, &LcEventSrc::Scanner).await; + let second = state.queue_transition_task(&other, &event, &LcEventSrc::Scanner).await; assert!(first); assert!(!second); @@ -3267,8 +3328,13 @@ mod tests { ..Default::default() }; + let other = ObjectInfo { + bucket: "bucket".to_string(), + name: "object-2".to_string(), + ..Default::default() + }; let first = state.queue_transition_task(&object, &event, &LcEventSrc::Scanner).await; - let second = state.queue_transition_task(&object, &event, &LcEventSrc::Scanner).await; + let second = state.queue_transition_task(&other, &event, &LcEventSrc::Scanner).await; assert!(first); assert!(!second); @@ -3769,6 +3835,62 @@ mod tests { assert_eq!(state.scanner_transition_state_update().compensation_pending, 1); } + #[test] + fn transition_reserve_dedupes_same_object_version() { + // Regression for rustfs/backlog#1268: a single object version can be + // enqueued twice (immediate transition + compensation backfill). The + // second claim must be rejected until the first is released, so the two + // do not run concurrently and race the source cleanup. + let state = TransitionState::new_with_capacity(4); + + let oi = ObjectInfo { + bucket: "foo".to_string(), + name: "payload.bin".to_string(), + version_id: None, + ..Default::default() + }; + + assert!(state.reserve_transition(&oi), "first claim must succeed"); + assert!(!state.reserve_transition(&oi), "duplicate claim must be rejected"); + + // A different object is independent. + let other = ObjectInfo { + bucket: "foo".to_string(), + name: "other.bin".to_string(), + version_id: None, + ..Default::default() + }; + assert!(state.reserve_transition(&other), "distinct object must claim independently"); + + // After release, the same object can be re-claimed (later lifecycle pass). + state.release_transition(&oi); + assert!(state.reserve_transition(&oi), "re-claim after release must succeed"); + } + + #[tokio::test] + #[serial] + async fn queue_transition_task_dedupes_same_object_without_second_enqueue() { + // Capacity 4 leaves room, so a rejected second enqueue can only be the + // dedup guard, not queue pressure (rustfs/backlog#1268). + let state = TransitionState::new_with_capacity(4); + let object = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + ..Default::default() + }; + let event = crate::bucket::lifecycle::lifecycle::Event { + action: IlmAction::TransitionAction, + ..Default::default() + }; + + let first = state.queue_transition_task(&object, &event, &LcEventSrc::Scanner).await; + let second = state.queue_transition_task(&object, &event, &LcEventSrc::Scanner).await; + + assert!(first, "first enqueue must succeed"); + assert!(second, "duplicate enqueue is reported handled (already queued)"); + assert_eq!(state.transition_rx.len(), 1, "only one task must actually be queued"); + } + #[tokio::test] #[serial] async fn transition_state_init_honors_runtime_configured_worker_count() {