diff --git a/crates/obs/src/metrics/collectors/replication.rs b/crates/obs/src/metrics/collectors/replication.rs index 15e1b7ae9..88b4b19a1 100644 --- a/crates/obs/src/metrics/collectors/replication.rs +++ b/crates/obs/src/metrics/collectors/replication.rs @@ -49,7 +49,7 @@ pub struct ReplicationStats { pub max_queued_count: u64, /// Maximum data transfer rate seen since server start pub max_data_transfer_rate: f64, - /// Objects in replication backlog in the last 5 minutes + /// Objects currently in replication backlog pub recent_backlog_count: u64, } diff --git a/crates/obs/src/metrics/schema/replication.rs b/crates/obs/src/metrics/schema/replication.rs index b9ee5e4df..6b857aa5b 100644 --- a/crates/obs/src/metrics/schema/replication.rs +++ b/crates/obs/src/metrics/schema/replication.rs @@ -128,7 +128,7 @@ pub static REPLICATION_MAX_DATA_TRANSFER_RATE_MD: LazyLock = L pub static REPLICATION_RECENT_BACKLOG_COUNT_MD: LazyLock = LazyLock::new(|| { new_gauge_md( MetricName::ReplicationRecentBacklogCount, - "Total number of objects seen in replication backlog in the last 5 minutes", + "Total number of objects currently in replication backlog (failed plus queued)", &[], subsystems::REPLICATION, ) diff --git a/crates/obs/src/metrics/storage_api.rs b/crates/obs/src/metrics/storage_api.rs index 040b34cee..a1fed693e 100644 --- a/crates/obs/src/metrics/storage_api.rs +++ b/crates/obs/src/metrics/storage_api.rs @@ -96,6 +96,12 @@ fn i32_to_u64_floor_zero(value: i32) -> u64 { u64::try_from(value.max(0)).unwrap_or(0) } +fn replication_backlog_count(failed_counts: impl Iterator, queued_count: i64) -> u64 { + let failed_backlog = failed_counts.map(i64_to_u64_floor_zero).sum::(); + + failed_backlog.saturating_add(i64_to_u64_floor_zero(queued_count)) +} + pub(crate) async fn obs_bucket_replication_stats_snapshot() -> Vec { let Some(stats) = get_global_replication_stats() else { return Vec::new(); @@ -191,7 +197,13 @@ pub(crate) async fn obs_replication_site_stats_snapshot(current_data_transfer_ra .flat_map(|bucket| bucket.stats.values()) .map(|stat| stat.xfer_rate_lrg.peak + stat.xfer_rate_sml.peak) .sum::(); - let recent_backlog_count = stats.mrf_stats.values().copied().filter(|value| *value > 0).sum::(); + let recent_backlog_count = replication_backlog_count( + all_bucket_stats + .values() + .flat_map(|bucket| bucket.stats.values()) + .map(|stat| stat.failed.count), + site_metrics.queued.curr.count, + ); ObsReplicationSiteStatsSnapshot { average_active_workers: site_metrics.active_workers.avg, @@ -206,7 +218,7 @@ pub(crate) async fn obs_replication_site_stats_snapshot(current_data_transfer_ra max_queued_bytes: i64_to_u64_floor_zero(site_metrics.queued.max.bytes), max_queued_count: i64_to_u64_floor_zero(site_metrics.queued.max.count), max_data_transfer_rate, - recent_backlog_count: i64_to_u64_floor_zero(recent_backlog_count), + recent_backlog_count, } } @@ -221,6 +233,16 @@ mod tests { assert_eq!(i32_to_u64_floor_zero(-1), 0); assert_eq!(i32_to_u64_floor_zero(42), 42); } + + #[test] + fn replication_backlog_count_uses_failed_targets_and_current_queue() { + assert_eq!(replication_backlog_count([3, 5].into_iter(), 7), 15); + } + + #[test] + fn replication_backlog_count_floors_negative_failed_and_queue_values() { + assert_eq!(replication_backlog_count([-3, 4].into_iter(), -2), 4); + } } pub(crate) mod metrics {