diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 903c1c24f..64e98b847 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -68,6 +68,30 @@ pub const MRF_WORKER_AUTO_DEFAULT: usize = 4; pub const LARGE_WORKER_COUNT: usize = 10; pub const MIN_LARGE_OBJ_SIZE: i64 = 128 * 1024 * 1024; // 128MiB +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub enum ReplicationQueueAdmission { + #[default] + Skipped, + Queued, + Missed, +} + +impl ReplicationQueueAdmission { + fn merge(&mut self, other: Self) { + *self = match (*self, other) { + (Self::Missed, _) | (_, Self::Missed) => Self::Missed, + (Self::Queued, _) | (_, Self::Queued) => Self::Queued, + (Self::Skipped, Self::Skipped) => Self::Skipped, + }; + } +} + +#[derive(Debug, Clone, Default)] +pub struct ReplicationHealQueueResult { + pub object_info: ReplicateObjectInfo, + pub admission: ReplicationQueueAdmission, +} + /// Priority levels for replication #[derive(Debug, Clone, PartialEq)] pub enum ReplicationPriority { @@ -513,7 +537,7 @@ impl ReplicationPool { } /// Queues a replica task - pub async fn queue_replica_task(&self, ri: ReplicateObjectInfo) { + pub async fn queue_replica_task(&self, ri: ReplicateObjectInfo) -> ReplicationQueueAdmission { // If object is large, queue it to a static set of large workers if ri.size >= MIN_LARGE_OBJ_SIZE { use std::collections::hash_map::DefaultHasher; @@ -532,7 +556,11 @@ impl ReplicationPool { && worker.try_send(ReplicationOperation::Object(Box::new(ri.clone()))).is_err() { // Queue to MRF if worker is busy - let _ = self.mrf_save_tx.try_send(ri.to_mrf_entry()); + let admission = if self.mrf_save_tx.try_send(ri.to_mrf_entry()).is_ok() { + ReplicationQueueAdmission::Queued + } else { + ReplicationQueueAdmission::Missed + }; // Try to add more workers if possible let max_l_workers = *self.max_l_workers.read().await; @@ -543,9 +571,12 @@ impl ReplicationPool { drop(lrg_workers); self.resize_lrg_workers(workers, existing).await; } + return admission; } + + return ReplicationQueueAdmission::Queued; } - return; + return ReplicationQueueAdmission::Missed; } // Handle regular sized objects @@ -555,85 +586,105 @@ impl ReplicationPool { _ => self.get_worker_ch(&ri.bucket, &ri.name, ri.size).await, }; - if let Some(channel) = ch - && channel.try_send(ReplicationOperation::Object(Box::new(ri.clone()))).is_err() - { - // Queue to MRF if all workers are busy - let _ = self.mrf_save_tx.try_send(ri.to_mrf_entry()); + let Some(channel) = ch else { + return ReplicationQueueAdmission::Missed; + }; - // Try to scale up workers based on priority - let priority = self.priority.read().await.clone(); - let max_workers = *self.max_workers.read().await; + if channel.try_send(ReplicationOperation::Object(Box::new(ri.clone()))).is_ok() { + return ReplicationQueueAdmission::Queued; + } - match priority { - ReplicationPriority::Fast => { - // Log warning about unable to keep up - info!("Warning: Unable to keep up with incoming traffic"); + // Queue to MRF if all workers are busy + let admission = if self.mrf_save_tx.try_send(ri.to_mrf_entry()).is_ok() { + ReplicationQueueAdmission::Queued + } else { + ReplicationQueueAdmission::Missed + }; + + // Try to scale up workers based on priority + let priority = self.priority.read().await.clone(); + let max_workers = *self.max_workers.read().await; + + match priority { + ReplicationPriority::Fast => { + // Log warning about unable to keep up + info!("Warning: Unable to keep up with incoming traffic"); + } + ReplicationPriority::Slow => { + info!("Warning: Unable to keep up with incoming traffic - recommend increasing replication priority to auto"); + } + ReplicationPriority::Auto => { + let max_w = std::cmp::min(max_workers, WORKER_MAX_LIMIT); + let active_workers = self.active_workers(); + + if active_workers < max_w as i32 { + let workers = self.workers.read().await; + let new_count = std::cmp::min(workers.len() + 1, max_w); + let existing = workers.len(); + + drop(workers); + self.resize_workers(new_count, existing).await; } - ReplicationPriority::Slow => { - info!("Warning: Unable to keep up with incoming traffic - recommend increasing replication priority to auto"); - } - ReplicationPriority::Auto => { - let max_w = std::cmp::min(max_workers, WORKER_MAX_LIMIT); - let active_workers = self.active_workers(); - if active_workers < max_w as i32 { - let workers = self.workers.read().await; - let new_count = std::cmp::min(workers.len() + 1, max_w); - let existing = workers.len(); + let max_mrf_workers = std::cmp::min(max_workers, MRF_WORKER_MAX_LIMIT); + let active_mrf = self.active_mrf_workers(); - drop(workers); - self.resize_workers(new_count, existing).await; - } + if active_mrf < max_mrf_workers as i32 { + let current_mrf = self.mrf_worker_size.load(Ordering::SeqCst); + let new_mrf = std::cmp::min(current_mrf + 1, max_mrf_workers as i32); - let max_mrf_workers = std::cmp::min(max_workers, MRF_WORKER_MAX_LIMIT); - let active_mrf = self.active_mrf_workers(); - - if active_mrf < max_mrf_workers as i32 { - let current_mrf = self.mrf_worker_size.load(Ordering::SeqCst); - let new_mrf = std::cmp::min(current_mrf + 1, max_mrf_workers as i32); - - self.resize_failed_workers(new_mrf).await; - } + self.resize_failed_workers(new_mrf).await; } } } + + admission } /// Queues a replica delete task - pub async fn queue_replica_delete_task(&self, doi: DeletedObjectReplicationInfo) { + pub async fn queue_replica_delete_task(&self, doi: DeletedObjectReplicationInfo) -> ReplicationQueueAdmission { let ch = match doi.op_type { ReplicationType::Heal | ReplicationType::ExistingObject => Some(self.mrf_replica_tx.clone()), _ => self.get_worker_ch(&doi.bucket, &doi.delete_object.object_name, 0).await, }; - if let Some(channel) = ch - && channel.try_send(ReplicationOperation::Delete(Box::new(doi.clone()))).is_err() - { - let _ = self.mrf_save_tx.try_send(doi.to_mrf_entry()); + let Some(channel) = ch else { + return ReplicationQueueAdmission::Missed; + }; - let priority = self.priority.read().await.clone(); - let max_workers = *self.max_workers.read().await; + if channel.try_send(ReplicationOperation::Delete(Box::new(doi.clone()))).is_ok() { + return ReplicationQueueAdmission::Queued; + } - match priority { - ReplicationPriority::Fast => { - info!("Warning: Unable to keep up with incoming deletes"); - } - ReplicationPriority::Slow => { - info!("Warning: Unable to keep up with incoming deletes - recommend increasing replication priority to auto"); - } - ReplicationPriority::Auto => { - let max_w = std::cmp::min(max_workers, WORKER_MAX_LIMIT); - if self.active_workers() < max_w as i32 { - let workers = self.workers.read().await; - let new_count = std::cmp::min(workers.len() + 1, max_w); - let existing = workers.len(); - drop(workers); - self.resize_workers(new_count, existing).await; - } + let admission = if self.mrf_save_tx.try_send(doi.to_mrf_entry()).is_ok() { + ReplicationQueueAdmission::Queued + } else { + ReplicationQueueAdmission::Missed + }; + + let priority = self.priority.read().await.clone(); + let max_workers = *self.max_workers.read().await; + + match priority { + ReplicationPriority::Fast => { + info!("Warning: Unable to keep up with incoming deletes"); + } + ReplicationPriority::Slow => { + info!("Warning: Unable to keep up with incoming deletes - recommend increasing replication priority to auto"); + } + ReplicationPriority::Auto => { + let max_w = std::cmp::min(max_workers, WORKER_MAX_LIMIT); + if self.active_workers() < max_w as i32 { + let workers = self.workers.read().await; + let new_count = std::cmp::min(workers.len() + 1, max_w); + let existing = workers.len(); + drop(workers); + self.resize_workers(new_count, existing).await; } } } + + admission } /// Queues an MRF save operation @@ -959,8 +1010,8 @@ pub trait ReplicationPoolTrait: std::fmt::Debug { fn active_workers(&self) -> i32; fn active_mrf_workers(&self) -> i32; fn active_lrg_workers(&self) -> i32; - async fn queue_replica_task(&self, ri: ReplicateObjectInfo); - async fn queue_replica_delete_task(&self, ri: DeletedObjectReplicationInfo); + async fn queue_replica_task(&self, ri: ReplicateObjectInfo) -> ReplicationQueueAdmission; + async fn queue_replica_delete_task(&self, ri: DeletedObjectReplicationInfo) -> ReplicationQueueAdmission; async fn resize(&self, priority: ReplicationPriority, max_workers: usize, max_l_workers: usize); async fn get_bucket_resync_status(&self, bucket: &str) -> Result; async fn cancel_bucket_resync(&self, opts: ResyncOpts) -> Result<(), EcstoreError>; @@ -987,12 +1038,12 @@ impl ReplicationPoolTrait for ReplicationPool { ReplicationPool::::active_lrg_workers(self) } - async fn queue_replica_task(&self, ri: ReplicateObjectInfo) { - self.queue_replica_task(ri).await; + async fn queue_replica_task(&self, ri: ReplicateObjectInfo) -> ReplicationQueueAdmission { + self.queue_replica_task(ri).await } - async fn queue_replica_delete_task(&self, ri: DeletedObjectReplicationInfo) { - self.queue_replica_delete_task(ri).await; + async fn queue_replica_delete_task(&self, ri: DeletedObjectReplicationInfo) -> ReplicationQueueAdmission { + self.queue_replica_delete_task(ri).await } async fn resize(&self, priority: ReplicationPriority, max_workers: usize, max_l_workers: usize) { @@ -1093,13 +1144,13 @@ pub async fn schedule_replication(oi: ObjectInfo, o: Arc, dsc: if dsc.is_synchronous() { replicate_object(ri, o).await } else if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { - pool.queue_replica_task(ri).await; + let _ = pool.queue_replica_task(ri).await; } } pub async fn schedule_replication_delete(dv: DeletedObjectReplicationInfo) { if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { - pool.queue_replica_delete_task(dv.clone()).await; + let _ = pool.queue_replica_delete_task(dv.clone()).await; } if let (Some(rs), Some(stats)) = (dv.delete_object.replication_state, GLOBAL_REPLICATION_STATS.get()) { @@ -1153,23 +1204,32 @@ pub async fn queue_replication_heal_internal( oi: ObjectInfo, rcfg: ReplicationConfig, retry_count: u32, -) -> ReplicateObjectInfo { +) -> ReplicationHealQueueResult { let mut roi = ReplicateObjectInfo::default(); // ignore modtime zero objects if oi.mod_time.is_none() || oi.mod_time == Some(OffsetDateTime::UNIX_EPOCH) { - return roi; + return ReplicationHealQueueResult { + object_info: roi, + admission: ReplicationQueueAdmission::Skipped, + }; } if rcfg.config.is_none() || rcfg.remotes.is_none() { - return roi; + return ReplicationHealQueueResult { + object_info: roi, + admission: ReplicationQueueAdmission::Skipped, + }; } roi = get_heal_replicate_object_info(&oi, &rcfg).await; roi.retry_count = retry_count; if !roi.dsc.replicate_any() { - return roi; + return ReplicationHealQueueResult { + object_info: roi, + admission: ReplicationQueueAdmission::Skipped, + }; } // early return if replication already done, otherwise we need to determine if this @@ -1178,7 +1238,10 @@ pub async fn queue_replication_heal_internal( && roi.version_purge_status.is_empty() && !roi.existing_obj_resync.must_resync() { - return roi; + return ReplicationHealQueueResult { + object_info: roi, + admission: ReplicationQueueAdmission::Skipped, + }; } if roi.delete_marker || !roi.version_purge_status.is_empty() { @@ -1210,10 +1273,15 @@ pub async fn queue_replication_heal_internal( || roi.version_purge_status == VersionPurgeStatusType::Failed || roi.version_purge_status == VersionPurgeStatusType::Pending { - if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { - pool.queue_replica_delete_task(dv).await; - } - return roi; + let admission = if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { + pool.queue_replica_delete_task(dv).await + } else { + ReplicationQueueAdmission::Missed + }; + return ReplicationHealQueueResult { + object_info: roi, + admission, + }; } // if replication status is Complete on DeleteMarker and existing object resync required @@ -1221,11 +1289,17 @@ pub async fn queue_replication_heal_internal( if existing_obj_resync.must_resync() && (roi.replication_status == ReplicationStatusType::Completed || roi.replication_status.is_empty()) { - queue_replicate_deletes_wrapper(dv, existing_obj_resync).await; - return roi; + let admission = queue_replicate_deletes_wrapper(dv, existing_obj_resync).await; + return ReplicationHealQueueResult { + object_info: roi, + admission, + }; } - return roi; + return ReplicationHealQueueResult { + object_info: roi, + admission: ReplicationQueueAdmission::Skipped, + }; } if roi.existing_obj_resync.must_resync() { @@ -1235,34 +1309,72 @@ pub async fn queue_replication_heal_internal( match roi.replication_status { ReplicationStatusType::Pending | ReplicationStatusType::Failed => { roi.event_type = REPLICATE_HEAL.to_string(); - if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { - pool.queue_replica_task(roi.clone()).await; - } - return roi; + let admission = if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { + pool.queue_replica_task(roi.clone()).await + } else { + ReplicationQueueAdmission::Missed + }; + return ReplicationHealQueueResult { + object_info: roi, + admission, + }; } _ => {} } if roi.existing_obj_resync.must_resync() { roi.event_type = REPLICATE_EXISTING.to_string(); - if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { - pool.queue_replica_task(roi.clone()).await; - } + let admission = if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { + pool.queue_replica_task(roi.clone()).await + } else { + ReplicationQueueAdmission::Missed + }; + return ReplicationHealQueueResult { + object_info: roi, + admission, + }; } - roi + ReplicationHealQueueResult { + object_info: roi, + admission: ReplicationQueueAdmission::Skipped, + } } /// Wrapper function for queueing replicate deletes with resync decision -async fn queue_replicate_deletes_wrapper(doi: DeletedObjectReplicationInfo, existing_obj_resync: ResyncDecision) { +async fn queue_replicate_deletes_wrapper( + doi: DeletedObjectReplicationInfo, + existing_obj_resync: ResyncDecision, +) -> ReplicationQueueAdmission { + let mut admission = ReplicationQueueAdmission::Skipped; for (k, v) in existing_obj_resync.targets.iter() { if v.replicate { let mut dv = doi.clone(); dv.reset_id = v.reset_id.clone(); dv.target_arn = k.clone(); - if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { - pool.queue_replica_delete_task(dv).await; - } + let target_admission = if let Some(pool) = GLOBAL_REPLICATION_POOL.get() { + pool.queue_replica_delete_task(dv).await + } else { + ReplicationQueueAdmission::Missed + }; + admission.merge(target_admission); } } + admission +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn replication_queue_admission_combines_target_results() { + let mut admission = ReplicationQueueAdmission::Skipped; + + admission.merge(ReplicationQueueAdmission::Queued); + assert_eq!(admission, ReplicationQueueAdmission::Queued); + + admission.merge(ReplicationQueueAdmission::Missed); + assert_eq!(admission, ReplicationQueueAdmission::Missed); + } } diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index bd53a0b2d..782527308 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -36,7 +36,8 @@ use rustfs_common::heal_channel::{ send_heal_request_with_admission, }; use rustfs_common::metrics::{ - IlmAction, Metric, Metrics, ScannerWorkSource, UpdateCurrentPathFn, current_path_updater, global_metrics, + IlmAction, Metric, Metrics, ScannerSourceWorkUpdate, ScannerWorkSource, UpdateCurrentPathFn, current_path_updater, + global_metrics, }; use rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc; use rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_ops::{GLOBAL_ExpiryState, apply_expiry_rule}; @@ -45,7 +46,9 @@ use rustfs_ecstore::bucket::lifecycle::{ bucket_lifecycle_ops::apply_transition_rule, lifecycle::{Event, Lifecycle, ObjectOpts}, }; -use rustfs_ecstore::bucket::replication::{ReplicationConfig, ReplicationConfigurationExt as _, queue_replication_heal_internal}; +use rustfs_ecstore::bucket::replication::{ + ReplicationConfig, ReplicationConfigurationExt as _, ReplicationQueueAdmission, queue_replication_heal_internal, +}; use rustfs_ecstore::bucket::versioning::VersioningApi; use rustfs_ecstore::bucket::versioning_sys::BucketVersioningSys; use rustfs_ecstore::cache_value::metacache_set::{ListPathRawOptions, list_path_raw}; @@ -161,6 +164,18 @@ fn record_scanner_ilm_action_if_queued(metrics: &Metrics, count: u64, queued: bo queued } +fn record_scanner_replication_admission(metrics: &Metrics, admission: ReplicationQueueAdmission) { + let work = match admission { + ReplicationQueueAdmission::Queued => ScannerSourceWorkUpdate::queued(1), + ReplicationQueueAdmission::Missed => ScannerSourceWorkUpdate::missed(1), + ReplicationQueueAdmission::Skipped => ScannerSourceWorkUpdate { + skipped: 1, + ..Default::default() + }, + }; + metrics.record_scanner_source_work(ScannerWorkSource::BucketReplication, work); +} + #[derive(Clone, Copy)] struct PendingScannerAccounting<'a> { object: &'a ObjectInfo, @@ -731,8 +746,10 @@ impl ScannerItem { }; let done_replication = Metrics::time(Metric::CheckReplication); - let roi = queue_replication_heal_internal(&oi.bucket, oi.clone(), (*replication).clone(), 0).await; + let replication_result = queue_replication_heal_internal(&oi.bucket, oi.clone(), (*replication).clone(), 0).await; done_replication(); + record_scanner_replication_admission(global_metrics(), replication_result.admission); + let roi = replication_result.object_info; if !Self::should_account_replication_stats(oi) { return; } @@ -2128,6 +2145,26 @@ mod tests { assert_eq!(queued_cumulative_size, 0); } + #[tokio::test] + async fn test_scanner_replication_admission_accounting_maps_source_work() { + let metrics = Metrics::new(); + + record_scanner_replication_admission(&metrics, ReplicationQueueAdmission::Skipped); + record_scanner_replication_admission(&metrics, ReplicationQueueAdmission::Queued); + record_scanner_replication_admission(&metrics, ReplicationQueueAdmission::Missed); + + let report = metrics.report().await; + let replication = report + .source_work + .iter() + .find(|work| work.source == ScannerWorkSource::BucketReplication.as_str()) + .expect("bucket replication source work should be visible"); + + assert_eq!(replication.skipped, 1); + assert_eq!(replication.queued, 1); + assert_eq!(replication.missed, 1); + } + #[test] #[serial] fn test_excessive_version_alert_thresholds_use_env() {