feat(replication): expose scanner repair outcomes (#3278)

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: cxymds <Cxymds@qq.com>
This commit is contained in:
Henry Guo
2026-06-08 22:24:27 +08:00
committed by GitHub
parent 452c003341
commit 6413df4b7f
2 changed files with 244 additions and 95 deletions
@@ -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<S: StorageAPI> ReplicationPool<S> {
}
/// 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<S: StorageAPI> ReplicationPool<S> {
&& 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<S: StorageAPI> ReplicationPool<S> {
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<S: StorageAPI> ReplicationPool<S> {
_ => 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<BucketReplicationResyncStatus, EcstoreError>;
async fn cancel_bucket_resync(&self, opts: ResyncOpts) -> Result<(), EcstoreError>;
@@ -987,12 +1038,12 @@ impl<S: StorageAPI> ReplicationPoolTrait for ReplicationPool<S> {
ReplicationPool::<S>::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<S: StorageAPI>(oi: ObjectInfo, o: Arc<S>, 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);
}
}
+40 -3
View File
@@ -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() {