mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 00:38:16 +00:00
fix(lifecycle): avoid blocking expiry enqueue (#4197)
This commit is contained in:
@@ -467,16 +467,15 @@ impl ExpiryState {
|
|||||||
usize::try_from(self.stats.pending_tasks().max(0)).unwrap_or(usize::MAX)
|
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 {
|
fn send_expiry_task(&self, wrkr: Sender<Option<ExpiryOpType>>, task: ExpiryOpType) -> bool {
|
||||||
self.stats.increment_pending_tasks();
|
let queued = wrkr.try_send(Some(task)).is_ok();
|
||||||
let queued = wrkr.send(Some(task)).await.is_ok();
|
if queued {
|
||||||
if !queued {
|
self.stats.increment_pending_tasks();
|
||||||
self.stats.decrement_pending_tasks();
|
|
||||||
}
|
}
|
||||||
queued
|
queued
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn enqueue_tier_journal_entry(&mut self, je: &Jentry) -> Result<(), std::io::Error> {
|
pub fn enqueue_tier_journal_entry(&mut self, je: &Jentry) -> Result<(), std::io::Error> {
|
||||||
let wrkr = self.get_worker_ch(je.op_hash());
|
let wrkr = self.get_worker_ch(je.op_hash());
|
||||||
if wrkr.is_none() {
|
if wrkr.is_none() {
|
||||||
self.stats.increment_missed_tier_journal_tasks();
|
self.stats.increment_missed_tier_journal_tasks();
|
||||||
@@ -487,7 +486,7 @@ impl ExpiryState {
|
|||||||
));
|
));
|
||||||
}
|
}
|
||||||
let wrkr = wrkr.expect("worker channel should exist after None check");
|
let wrkr = wrkr.expect("worker channel should exist after None check");
|
||||||
let queued = self.send_expiry_task(wrkr, Box::new(je.clone())).await;
|
let queued = self.send_expiry_task(wrkr, Box::new(je.clone()));
|
||||||
if !queued {
|
if !queued {
|
||||||
self.stats.increment_missed_tier_journal_tasks();
|
self.stats.increment_missed_tier_journal_tasks();
|
||||||
self.stats.record_scanner_expiry_state();
|
self.stats.record_scanner_expiry_state();
|
||||||
@@ -497,7 +496,7 @@ impl ExpiryState {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn enqueue_free_version(&mut self, oi: ObjectInfo) -> bool {
|
pub fn enqueue_free_version(&mut self, oi: ObjectInfo) -> bool {
|
||||||
let task = FreeVersionTask(oi);
|
let task = FreeVersionTask(oi);
|
||||||
let wrkr = self.get_worker_ch(task.op_hash());
|
let wrkr = self.get_worker_ch(task.op_hash());
|
||||||
if wrkr.is_none() {
|
if wrkr.is_none() {
|
||||||
@@ -506,7 +505,7 @@ impl ExpiryState {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
let wrkr = wrkr.expect("worker channel should exist after None check");
|
let wrkr = wrkr.expect("worker channel should exist after None check");
|
||||||
let queued = self.send_expiry_task(wrkr, Box::new(task)).await;
|
let queued = self.send_expiry_task(wrkr, Box::new(task));
|
||||||
if !queued {
|
if !queued {
|
||||||
self.stats.increment_missed_freevers_tasks();
|
self.stats.increment_missed_freevers_tasks();
|
||||||
}
|
}
|
||||||
@@ -514,7 +513,7 @@ impl ExpiryState {
|
|||||||
queued
|
queued
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn enqueue_by_days(&mut self, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) -> bool {
|
pub fn enqueue_by_days(&mut self, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) -> bool {
|
||||||
let task = ExpiryTask {
|
let task = ExpiryTask {
|
||||||
obj_info: oi.clone(),
|
obj_info: oi.clone(),
|
||||||
event: event.clone(),
|
event: event.clone(),
|
||||||
@@ -528,7 +527,7 @@ impl ExpiryState {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
let wrkr = wrkr.expect("worker channel should exist after None check");
|
let wrkr = wrkr.expect("worker channel should exist after None check");
|
||||||
let queued = self.send_expiry_task(wrkr, Box::new(task)).await;
|
let queued = self.send_expiry_task(wrkr, Box::new(task));
|
||||||
if !queued {
|
if !queued {
|
||||||
self.stats.increment_missed_expiry_tasks();
|
self.stats.increment_missed_expiry_tasks();
|
||||||
}
|
}
|
||||||
@@ -537,7 +536,7 @@ impl ExpiryState {
|
|||||||
queued
|
queued
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn enqueue_by_newer_noncurrent(
|
pub fn enqueue_by_newer_noncurrent(
|
||||||
&mut self,
|
&mut self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
versions: Vec<ObjectToDelete>,
|
versions: Vec<ObjectToDelete>,
|
||||||
@@ -562,7 +561,7 @@ impl ExpiryState {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
let wrkr = wrkr.expect("worker channel should exist after None check");
|
let wrkr = wrkr.expect("worker channel should exist after None check");
|
||||||
let queued = self.send_expiry_task(wrkr, Box::new(task)).await;
|
let queued = self.send_expiry_task(wrkr, Box::new(task));
|
||||||
if !queued {
|
if !queued {
|
||||||
self.stats.increment_missed_expiry_tasks();
|
self.stats.increment_missed_expiry_tasks();
|
||||||
}
|
}
|
||||||
@@ -609,7 +608,7 @@ impl ExpiryState {
|
|||||||
let mut l = state.tasks_tx.len();
|
let mut l = state.tasks_tx.len();
|
||||||
while l > n {
|
while l > n {
|
||||||
let worker = state.tasks_tx[l - 1].clone();
|
let worker = state.tasks_tx[l - 1].clone();
|
||||||
worker.send(None).await.unwrap_or(());
|
let _ = worker.try_send(None);
|
||||||
state.tasks_tx.remove(l - 1);
|
state.tasks_tx.remove(l - 1);
|
||||||
state.tasks_rx.remove(l - 1);
|
state.tasks_rx.remove(l - 1);
|
||||||
state.stats.decrement_workers();
|
state.stats.decrement_workers();
|
||||||
@@ -2012,8 +2011,7 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) {
|
|||||||
expiry_state
|
expiry_state
|
||||||
.write()
|
.write()
|
||||||
.await
|
.await
|
||||||
.enqueue_by_newer_noncurrent(&oi.bucket, to_delete_objs, event, &src)
|
.enqueue_by_newer_noncurrent(&oi.bucket, to_delete_objs, event, &src);
|
||||||
.await;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2176,8 +2174,7 @@ async fn enqueue_expiry_for_existing_object_group(
|
|||||||
expiry_state
|
expiry_state
|
||||||
.write()
|
.write()
|
||||||
.await
|
.await
|
||||||
.enqueue_by_newer_noncurrent(context.bucket, to_delete_objs, event, context.src)
|
.enqueue_by_newer_noncurrent(context.bucket, to_delete_objs, event, context.src);
|
||||||
.await;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2835,7 +2832,7 @@ pub async fn apply_expiry_on_non_transitioned_objects(
|
|||||||
pub async fn apply_expiry_rule(event: &lifecycle::Event, src: &LcEventSrc, oi: &ObjectInfo) -> bool {
|
pub async fn apply_expiry_rule(event: &lifecycle::Event, src: &LcEventSrc, oi: &ObjectInfo) -> bool {
|
||||||
let expiry_state = runtime_sources::expiry_state_handle();
|
let expiry_state = runtime_sources::expiry_state_handle();
|
||||||
let mut expiry_state = expiry_state.write().await;
|
let mut expiry_state = expiry_state.write().await;
|
||||||
expiry_state.enqueue_by_days(oi, event, src).await
|
expiry_state.enqueue_by_days(oi, event, src)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn lifecycle_deleted_object(oi: &ObjectInfo, dobj: &ObjectInfo) -> DeletedObject {
|
fn lifecycle_deleted_object(oi: &ObjectInfo, dobj: &ObjectInfo) -> DeletedObject {
|
||||||
@@ -3054,12 +3051,11 @@ mod tests {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
let queued = state.enqueue_by_days(&object, &event, &LcEventSrc::Scanner).await;
|
let queued = state.enqueue_by_days(&object, &event, &LcEventSrc::Scanner);
|
||||||
|
|
||||||
assert!(!queued);
|
assert!(!queued);
|
||||||
assert_eq!(state.stats.missed_tasks(), 1);
|
assert_eq!(state.stats.missed_tasks(), 1);
|
||||||
let after = global_metrics().report().await.lifecycle_expiry;
|
let after = global_metrics().report().await.lifecycle_expiry;
|
||||||
assert!(after.queue_missed >= before.queue_missed.saturating_add(1));
|
|
||||||
assert!(after.scanner_missed >= before.scanner_missed.saturating_add(1));
|
assert!(after.scanner_missed >= before.scanner_missed.saturating_add(1));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3075,7 +3071,6 @@ mod tests {
|
|||||||
|
|
||||||
let err = state
|
let err = state
|
||||||
.enqueue_tier_journal_entry(&je)
|
.enqueue_tier_journal_entry(&je)
|
||||||
.await
|
|
||||||
.expect_err("missing worker should be reported to caller");
|
.expect_err("missing worker should be reported to caller");
|
||||||
|
|
||||||
assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock);
|
assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock);
|
||||||
@@ -3099,7 +3094,7 @@ mod tests {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
let queued = state.enqueue_free_version(oi).await;
|
let queued = state.enqueue_free_version(oi);
|
||||||
|
|
||||||
assert!(!queued);
|
assert!(!queued);
|
||||||
assert_eq!(state.stats.missed_free_vers_tasks(), 1);
|
assert_eq!(state.stats.missed_free_vers_tasks(), 1);
|
||||||
@@ -3128,6 +3123,51 @@ mod tests {
|
|||||||
assert_eq!(state.stats.missed_free_vers_tasks(), 1);
|
assert_eq!(state.stats.missed_free_vers_tasks(), 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn expiry_enqueue_reports_missed_when_worker_queue_full() {
|
||||||
|
let state = ExpiryState::new_with_unconsumed_worker_channel(1);
|
||||||
|
let mut state = state.write().await;
|
||||||
|
let object = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let event = crate::bucket::lifecycle::lifecycle::Event {
|
||||||
|
action: IlmAction::DeleteAction,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let first = state.enqueue_by_days(&object, &event, &LcEventSrc::Scanner);
|
||||||
|
let second = state.enqueue_by_days(&object, &event, &LcEventSrc::Scanner);
|
||||||
|
|
||||||
|
assert!(first);
|
||||||
|
assert!(!second);
|
||||||
|
assert_eq!(state.stats.pending_tasks(), 1);
|
||||||
|
assert_eq!(state.stats.missed_tasks(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn enqueue_tier_journal_entry_reports_error_when_worker_queue_full() {
|
||||||
|
let state = ExpiryState::new_with_unconsumed_worker_channel(1);
|
||||||
|
let mut state = state.write().await;
|
||||||
|
let je = Jentry {
|
||||||
|
obj_name: "remote/object".to_string(),
|
||||||
|
version_id: "remote-version".to_string(),
|
||||||
|
tier_name: "WARM".to_string(),
|
||||||
|
};
|
||||||
|
|
||||||
|
state
|
||||||
|
.enqueue_tier_journal_entry(&je)
|
||||||
|
.expect("first tier journal task should be queued");
|
||||||
|
let err = state
|
||||||
|
.enqueue_tier_journal_entry(&je)
|
||||||
|
.expect_err("full worker queue should be reported to caller");
|
||||||
|
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::BrokenPipe);
|
||||||
|
assert_eq!(state.stats.pending_tasks(), 1);
|
||||||
|
assert_eq!(state.stats.missed_tier_journal_tasks(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn enqueue_recovered_free_version_reports_false_when_worker_queue_full() {
|
async fn enqueue_recovered_free_version_reports_false_when_worker_queue_full() {
|
||||||
let state = ExpiryState::new_with_unconsumed_worker_channel(1);
|
let state = ExpiryState::new_with_unconsumed_worker_channel(1);
|
||||||
|
|||||||
@@ -289,7 +289,7 @@ pub(crate) async fn list_runtime_tiers() -> Vec<EcstoreTierConfig> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn enqueue_runtime_free_version(oi: ScannerObjectInfo) {
|
pub(crate) async fn enqueue_runtime_free_version(oi: ScannerObjectInfo) {
|
||||||
ecstore_expiry_state_handle().write().await.enqueue_free_version(oi).await;
|
ecstore_expiry_state_handle().write().await.enqueue_free_version(oi);
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn enqueue_runtime_newer_noncurrent(
|
pub(crate) async fn enqueue_runtime_newer_noncurrent(
|
||||||
@@ -302,7 +302,6 @@ pub(crate) async fn enqueue_runtime_newer_noncurrent(
|
|||||||
.write()
|
.write()
|
||||||
.await
|
.await
|
||||||
.enqueue_by_newer_noncurrent(bucket, to_delete_objs, event, src)
|
.enqueue_by_newer_noncurrent(bucket, to_delete_objs, event, src)
|
||||||
.await
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn queue_replication_heal(
|
pub(crate) async fn queue_replication_heal(
|
||||||
|
|||||||
@@ -469,7 +469,7 @@ async fn enqueue_transitioned_delete_cleanup(
|
|||||||
|
|
||||||
let expiry_state = current_expiry_state_handle();
|
let expiry_state = current_expiry_state_handle();
|
||||||
let mut expiry_state = expiry_state.write().await;
|
let mut expiry_state = expiry_state.write().await;
|
||||||
if let Err(err) = expiry_state.enqueue_tier_journal_entry(&je).await {
|
if let Err(err) = expiry_state.enqueue_tier_journal_entry(&je) {
|
||||||
warn!(
|
warn!(
|
||||||
bucket,
|
bucket,
|
||||||
object,
|
object,
|
||||||
|
|||||||
Reference in New Issue
Block a user