mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-30 16:59:52 +00:00
fix(ecstore): deduplicate concurrent ILM transition enqueues (#4839)
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 <heihutu@gmail.com>
This commit is contained in:
@@ -836,6 +836,12 @@ pub struct TransitionState {
|
|||||||
compensation_scheduled_tasks: AtomicI64,
|
compensation_scheduled_tasks: AtomicI64,
|
||||||
compensation_running_tasks: AtomicI64,
|
compensation_running_tasks: AtomicI64,
|
||||||
compensation_buckets: Arc<Mutex<HashSet<String>>>,
|
compensation_buckets: Arc<Mutex<HashSet<String>>>,
|
||||||
|
// (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<Mutex<HashSet<(String, String, String)>>>,
|
||||||
last_day_stats: Arc<Mutex<HashMap<String, LastDayTierStats>>>,
|
last_day_stats: Arc<Mutex<HashMap<String, LastDayTierStats>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -869,10 +875,39 @@ impl TransitionState {
|
|||||||
compensation_scheduled_tasks: AtomicI64::new(0),
|
compensation_scheduled_tasks: AtomicI64::new(0),
|
||||||
compensation_running_tasks: AtomicI64::new(0),
|
compensation_running_tasks: AtomicI64::new(0),
|
||||||
compensation_buckets: Arc::new(Mutex::new(HashSet::new())),
|
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())),
|
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 {
|
fn reserve_bucket_compensation(&self, bucket: &str) -> bool {
|
||||||
let inserted = match self.compensation_buckets.lock() {
|
let inserted = match self.compensation_buckets.lock() {
|
||||||
Ok(mut scheduled) => scheduled.insert(bucket.to_string()),
|
Ok(mut scheduled) => scheduled.insert(bucket.to_string()),
|
||||||
@@ -1046,6 +1081,15 @@ impl TransitionState {
|
|||||||
return false;
|
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 {
|
let task = TransitionTask {
|
||||||
obj_info: oi.clone(),
|
obj_info: oi.clone(),
|
||||||
src: src.clone(),
|
src: src.clone(),
|
||||||
@@ -1084,6 +1128,9 @@ impl TransitionState {
|
|||||||
self.handle_immediate_enqueue_failure(oi, src, ImmediateEnqueueFailure::QueueClosed { timeout_ms: None });
|
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);
|
record_scanner_transition_enqueue_result(src, 1, queued);
|
||||||
self.record_scanner_transition_state();
|
self.record_scanner_transition_state();
|
||||||
return queued;
|
return queued;
|
||||||
@@ -1113,6 +1160,9 @@ impl TransitionState {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if !queued {
|
||||||
|
self.release_transition(oi);
|
||||||
|
}
|
||||||
record_scanner_transition_enqueue_result(src, 1, queued);
|
record_scanner_transition_enqueue_result(src, 1, queued);
|
||||||
self.record_scanner_transition_state();
|
self.record_scanner_transition_state();
|
||||||
queued
|
queued
|
||||||
@@ -1239,6 +1289,9 @@ impl TransitionState {
|
|||||||
|
|
||||||
emit_transition_complete_event(obj_info_for_event);
|
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);
|
TransitionState::add_counter(&transition_state.active_tasks, -1);
|
||||||
transition_state.record_scanner_transition_state();
|
transition_state.record_scanner_transition_state();
|
||||||
}
|
}
|
||||||
@@ -3244,8 +3297,16 @@ mod tests {
|
|||||||
..Default::default()
|
..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 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!(first);
|
||||||
assert!(!second);
|
assert!(!second);
|
||||||
@@ -3267,8 +3328,13 @@ mod tests {
|
|||||||
..Default::default()
|
..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 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!(first);
|
||||||
assert!(!second);
|
assert!(!second);
|
||||||
@@ -3769,6 +3835,62 @@ mod tests {
|
|||||||
assert_eq!(state.scanner_transition_state_update().compensation_pending, 1);
|
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]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn transition_state_init_honors_runtime_configured_worker_count() {
|
async fn transition_state_init_honors_runtime_configured_worker_count() {
|
||||||
|
|||||||
Reference in New Issue
Block a user