feat(obs): improve telemetry stack, replication metrics, and Grafana alignment (#2672)

Co-authored-by: Filipe Monteiro <a22407332@alunos.ulht.pt>
Co-authored-by: cxymds <Cxymds@qq.com>
Co-authored-by: weisd <im@weisd.in>
Co-authored-by: loverustfs <hello@rustfs.com>
Co-authored-by: 安正超 <anzhengchao@gmail.com>
This commit is contained in:
houseme
2026-04-24 21:50:17 +08:00
committed by GitHub
parent 2705e3f53b
commit 13b4500212
19 changed files with 4603 additions and 531 deletions
@@ -956,6 +956,9 @@ pub type DynReplicationPool = dyn ReplicationPoolTrait + Send + Sync;
/// Trait that abstracts the replication pool operations
#[async_trait::async_trait]
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 resize(&self, priority: ReplicationPriority, max_workers: usize, max_l_workers: usize);
@@ -972,6 +975,18 @@ pub trait ReplicationPoolTrait: std::fmt::Debug {
// Implement the trait for ReplicationPool
#[async_trait::async_trait]
impl<S: StorageAPI> ReplicationPoolTrait for ReplicationPool<S> {
fn active_workers(&self) -> i32 {
ReplicationPool::<S>::active_workers(self)
}
fn active_mrf_workers(&self) -> i32 {
ReplicationPool::<S>::active_mrf_workers(self)
}
fn active_lrg_workers(&self) -> i32 {
ReplicationPool::<S>::active_lrg_workers(self)
}
async fn queue_replica_task(&self, ri: ReplicateObjectInfo) {
self.queue_replica_task(ri).await;
}
@@ -35,7 +35,7 @@ use crate::set_disk::get_lock_acquire_timeout;
use crate::store_api::{DeletedObject, HTTPRangeSpec, ObjectInfo, ObjectOptions, ObjectToDelete, WalkOptions};
use crate::{StorageAPI, new_object_layer_fn};
use aws_sdk_s3::error::{ProvideErrorMetadata, SdkError};
use aws_sdk_s3::operation::head_object::HeadObjectOutput;
use aws_sdk_s3::operation::head_object::{HeadObjectError, HeadObjectOutput};
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{CompletedPart, ObjectLockLegalHoldStatus};
use aws_smithy_types::body::SdkBody;
@@ -115,6 +115,41 @@ fn resync_state_accepts_update(state: &TargetReplicationResyncStatus, opts: &Res
state.resync_id.is_empty() || opts.resync_id.is_empty() || state.resync_id == opts.resync_id
}
fn should_count_head_proxy_failure(is_not_found: bool, code: Option<&str>, raw_status: Option<u16>) -> bool {
if is_not_found || matches!(code, Some("MethodNotAllowed" | "405")) {
return false;
}
!matches!(raw_status, Some(404 | 405))
}
fn is_head_proxy_failure(err: &SdkError<HeadObjectError>) -> bool {
let (is_not_found, code) = err
.as_service_error()
.map(|service_err| (service_err.is_not_found(), service_err.code()))
.unwrap_or((false, None));
let raw_status = err.raw_response().map(|resp| resp.status().as_u16());
should_count_head_proxy_failure(is_not_found, code, raw_status)
}
async fn record_proxy_request(bucket: &str, api: &str, is_err: bool) {
if let Some(stats) = GLOBAL_REPLICATION_STATS.get() {
stats.inc_proxy(bucket, api, is_err).await;
}
}
async fn head_object_with_proxy_stats(
source_bucket: &str,
target_client: &TargetClient,
target_bucket: &str,
object: &str,
version_id: Option<String>,
) -> std::result::Result<HeadObjectOutput, SdkError<HeadObjectError>> {
let result = target_client.head_object(target_bucket, object, version_id).await;
let is_err = result.as_ref().err().is_some_and(is_head_proxy_failure);
record_proxy_request(source_bucket, "HeadObject", is_err).await;
result
}
#[derive(Debug, Clone, Default)]
pub struct ResyncOpts {
pub bucket: String,
@@ -748,10 +783,15 @@ impl ReplicationResyncer {
let reset_id = target_client.reset_id.clone();
let (size, err) = if let Err(err) = target_client
.head_object(&target_client.bucket, &roi.name, roi.version_id.map(|v| v.to_string()))
.await
{
let head_result = head_object_with_proxy_stats(
&bucket_name,
target_client.as_ref(),
&target_client.bucket,
&roi.name,
roi.version_id.map(|v| v.to_string()),
)
.await;
let (size, err) = if let Err(err) = head_result {
if roi.delete_marker {
st.replicated_count += 1;
} else {
@@ -2120,9 +2160,14 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli
};
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
match tgt_client
.head_object(&tgt_client.bucket, &dobj.delete_object.object_name, version_id.clone())
.await
match head_object_with_proxy_stats(
&dobj.bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
&dobj.delete_object.object_name,
version_id.clone(),
)
.await
{
Ok(_) => {}
Err(e) => {
@@ -2471,9 +2516,14 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
}
let mut replication_action = replication_action;
match tgt_client
.head_object(&tgt_client.bucket, &object, self.version_id.map(|v| v.to_string()))
.await
match head_object_with_proxy_stats(
&bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
&object,
self.version_id.map(|v| v.to_string()),
)
.await
{
Ok(oi) => {
replication_action = get_replication_action(&object_info, &oi, self.op_type);
@@ -2530,7 +2580,7 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
if let Some(err) = if is_multipart {
drop(gr);
replicate_object_with_multipart(MultipartReplicationContext {
let result = replicate_object_with_multipart(MultipartReplicationContext {
storage: storage.clone(),
cli: tgt_client.clone(),
src_bucket: &bucket,
@@ -2541,16 +2591,18 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
arn: &rinfo.arn,
put_opts,
})
.await
.err()
.await;
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
result.err()
} else {
gr.stream = wrap_with_bandwidth_monitor(gr.stream, &put_opts, &bucket, &rinfo.arn);
let byte_stream = async_read_to_bytestream(gr.stream);
tgt_client
let result = tgt_client
.put_object(&tgt_client.bucket, &object, size, byte_stream, &put_opts)
.await
.map_err(|e| std::io::Error::other(e.to_string()))
.err()
.map_err(|e| std::io::Error::other(e.to_string()));
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
result.err()
} {
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(err.to_string());
@@ -2690,9 +2742,14 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
warn!("failed to set replication tagging directive header: {err}");
}
match tgt_client
.head_object(&tgt_client.bucket, &object, self.version_id.map(|v| v.to_string()))
.await
match head_object_with_proxy_stats(
&bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
&object,
self.version_id.map(|v| v.to_string()),
)
.await
{
Ok(oi) => {
replication_action = get_replication_action(&object_info, &oi, self.op_type);
@@ -2812,7 +2869,7 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
if let Some(err) = if is_multipart {
drop(gr);
replicate_object_with_multipart(MultipartReplicationContext {
let result = replicate_object_with_multipart(MultipartReplicationContext {
storage: storage.clone(),
cli: tgt_client.clone(),
src_bucket: &bucket,
@@ -2823,16 +2880,18 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
arn: &rinfo.arn,
put_opts,
})
.await
.err()
.await;
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
result.err()
} else {
gr.stream = wrap_with_bandwidth_monitor(gr.stream, &put_opts, &bucket, &rinfo.arn);
let byte_stream = async_read_to_bytestream(gr.stream);
tgt_client
let result = tgt_client
.put_object(&tgt_client.bucket, &object, size, byte_stream, &put_opts)
.await
.map_err(|e| std::io::Error::other(e.to_string()))
.err()
.map_err(|e| std::io::Error::other(e.to_string()));
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
result.err()
} {
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(err.to_string());
@@ -3748,6 +3807,34 @@ mod tests {
);
}
#[test]
fn test_should_count_head_proxy_failure_ignores_not_found_and_405() {
assert!(
!should_count_head_proxy_failure(true, Some("NoSuchKey"), Some(404)),
"not-found heads are expected when the object has not reached the target yet"
);
assert!(
!should_count_head_proxy_failure(false, Some("MethodNotAllowed"), Some(405)),
"405 delete-marker probing responses should not be counted as proxy failures"
);
assert!(
!should_count_head_proxy_failure(false, Some("405"), Some(405)),
"numeric 405 codes must align with MethodNotAllowed semantics"
);
}
#[test]
fn test_should_count_head_proxy_failure_counts_unexpected_errors() {
assert!(
should_count_head_proxy_failure(false, Some("AccessDenied"), Some(403)),
"non-NotFound and non-405 service errors should be counted as failures"
);
assert!(
should_count_head_proxy_failure(false, None, Some(500)),
"raw 5xx head responses should be counted as proxy failures"
);
}
#[tokio::test]
async fn test_get_heal_replicate_object_info_failed_object_returns_heal_roi() {
let oi = ObjectInfo {
@@ -12,18 +12,22 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::bucket::replication::get_global_replication_pool;
use crate::error::Error;
use crate::global::get_global_bucket_monitor;
use rustfs_filemeta::{ReplicatedTargetInfo, ReplicationStatusType, ReplicationType};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::sync::atomic::{AtomicI64, Ordering};
use std::sync::atomic::{AtomicU64, Ordering as AtomicOrdering};
use std::time::{Duration, SystemTime};
use std::time::{Duration, Instant, SystemTime};
use tokio::sync::{Mutex, RwLock};
use tokio::time::interval;
const ROLLING_WINDOW: Duration = Duration::from_secs(60);
const FAILURE_LAST_HOUR_WINDOW: Duration = Duration::from_secs(60 * 60);
/// Exponential Moving Average with thread-safe interior mutability
#[derive(Debug)]
pub struct ExponentialMovingAverage {
@@ -328,6 +332,13 @@ pub struct InQueueStats {
pub now_count: AtomicI64,
}
#[derive(Debug, Clone)]
struct QueueSample {
observed_at: Instant,
bytes: i64,
count: i64,
}
impl Clone for InQueueStats {
fn clone(&self) -> Self {
Self {
@@ -359,9 +370,60 @@ pub struct InQueueMetric {
pub curr: InQueueStats,
pub avg: InQueueStats,
pub max: InQueueStats,
pub last_minute: InQueueStats,
#[serde(skip)]
samples: VecDeque<QueueSample>,
}
impl InQueueMetric {
fn observe(&mut self, observed_at: Instant) {
let bytes = self.curr.now_bytes.load(Ordering::Relaxed);
let count = self.curr.now_count.load(Ordering::Relaxed);
self.curr.bytes = bytes;
self.curr.count = count;
self.samples.push_back(QueueSample {
observed_at,
bytes,
count,
});
while self
.samples
.front()
.is_some_and(|sample| observed_at.duration_since(sample.observed_at) > ROLLING_WINDOW)
{
self.samples.pop_front();
}
if self.samples.is_empty() {
self.avg = InQueueStats::default();
self.max = InQueueStats::default();
self.last_minute = InQueueStats::default();
return;
}
let sample_count = self.samples.len() as i64;
let total_bytes = self.samples.iter().map(|sample| sample.bytes).sum::<i64>();
let total_count = self.samples.iter().map(|sample| sample.count).sum::<i64>();
let max_bytes = self.samples.iter().map(|sample| sample.bytes).max().unwrap_or(0);
let max_count = self.samples.iter().map(|sample| sample.count).max().unwrap_or(0);
self.avg.bytes = total_bytes / sample_count;
self.avg.count = total_count / sample_count;
self.max.bytes = max_bytes;
self.max.count = max_count;
self.last_minute.bytes = self.avg.bytes;
self.last_minute.count = self.avg.count;
}
fn snapshot(&self) -> Self {
let mut snapshot = self.clone();
snapshot.curr.bytes = snapshot.curr.now_bytes.load(Ordering::Relaxed);
snapshot.curr.count = snapshot.curr.now_count.load(Ordering::Relaxed);
snapshot
}
pub fn merge(&self, other: &InQueueMetric) -> Self {
Self {
curr: InQueueStats {
@@ -384,6 +446,12 @@ impl InQueueMetric {
count: self.max.count.max(other.max.count),
..Default::default()
},
last_minute: InQueueStats {
bytes: self.last_minute.bytes + other.last_minute.bytes,
count: self.last_minute.count + other.last_minute.count,
..Default::default()
},
samples: VecDeque::new(),
}
}
}
@@ -391,8 +459,8 @@ impl InQueueMetric {
/// Queue cache
#[derive(Debug, Default)]
pub struct QueueCache {
pub bucket_stats: HashMap<String, InQueueStats>,
pub sr_queue_stats: InQueueStats,
pub bucket_stats: HashMap<String, InQueueMetric>,
pub sr_queue_stats: InQueueMetric,
}
impl QueueCache {
@@ -401,36 +469,19 @@ impl QueueCache {
}
pub fn update(&mut self) {
// Update queue statistics cache
// In actual implementation, this would get latest statistics from queue system
let observed_at = Instant::now();
self.sr_queue_stats.observe(observed_at);
for stats in self.bucket_stats.values_mut() {
stats.observe(observed_at);
}
}
pub fn get_bucket_stats(&self, bucket: &str) -> InQueueMetric {
if let Some(bucket_stat) = self.bucket_stats.get(bucket) {
InQueueMetric {
curr: InQueueStats {
bytes: bucket_stat.now_bytes.load(Ordering::Relaxed),
count: bucket_stat.now_count.load(Ordering::Relaxed),
..Default::default()
},
avg: InQueueStats::default(), // simplified implementation
max: InQueueStats::default(), // simplified implementation
}
} else {
InQueueMetric::default()
}
self.bucket_stats.get(bucket).map(InQueueMetric::snapshot).unwrap_or_default()
}
pub fn get_site_stats(&self) -> InQueueMetric {
InQueueMetric {
curr: InQueueStats {
bytes: self.sr_queue_stats.now_bytes.load(Ordering::Relaxed),
count: self.sr_queue_stats.now_count.load(Ordering::Relaxed),
..Default::default()
},
avg: InQueueStats::default(), // simplified implementation
max: InQueueStats::default(), // simplified implementation
}
self.sr_queue_stats.snapshot()
}
}
@@ -505,11 +556,19 @@ impl ProxyStatsCache {
}
}
#[derive(Debug, Clone)]
struct FailureSample {
observed_at: Instant,
size: i64,
}
/// Failure statistics
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct FailStats {
pub count: i64,
pub size: i64,
#[serde(skip)]
recent: VecDeque<FailureSample>,
}
impl FailStats {
@@ -518,14 +577,42 @@ impl FailStats {
}
pub fn add_size(&mut self, size: i64, _err: Option<&Error>) {
let observed_at = Instant::now();
self.count += 1;
self.size += size;
self.recent.push_back(FailureSample { observed_at, size });
self.prune(observed_at);
}
fn prune(&mut self, observed_at: Instant) {
while self
.recent
.front()
.is_some_and(|sample| observed_at.duration_since(sample.observed_at) > FAILURE_LAST_HOUR_WINDOW)
{
self.recent.pop_front();
}
}
pub fn recent_since(&self, window: Duration) -> FailedMetric {
let now = Instant::now();
let mut count = 0i64;
let mut size = 0i64;
for sample in self.recent.iter().rev() {
if now.duration_since(sample.observed_at) > window {
break;
}
count += 1;
size += sample.size;
}
FailedMetric { count, size }
}
pub fn merge(&self, other: &FailStats) -> Self {
Self {
count: self.count + other.count,
size: self.size + other.size,
recent: VecDeque::new(),
}
}
@@ -674,6 +761,14 @@ pub struct ActiveWorkerStat {
pub curr: i32,
pub max: i32,
pub avg: f64,
#[serde(skip)]
samples: VecDeque<WorkerSample>,
}
#[derive(Debug, Clone)]
struct WorkerSample {
observed_at: Instant,
workers: i32,
}
impl ActiveWorkerStat {
@@ -685,9 +780,31 @@ impl ActiveWorkerStat {
self.clone()
}
pub fn update(&mut self) {
// Simulate worker statistics update logic
// In actual implementation, this would get current active count from worker pool
pub fn update(&mut self, curr: i32) {
let observed_at = Instant::now();
self.curr = curr;
self.samples.push_back(WorkerSample {
observed_at,
workers: curr,
});
while self
.samples
.front()
.is_some_and(|sample| observed_at.duration_since(sample.observed_at) > ROLLING_WINDOW)
{
self.samples.pop_front();
}
if self.samples.is_empty() {
self.max = curr;
self.avg = curr as f64;
return;
}
self.max = self.samples.iter().map(|sample| sample.workers).max().unwrap_or(curr);
let total = self.samples.iter().map(|sample| sample.workers as i64).sum::<i64>();
self.avg = total as f64 / self.samples.len() as f64;
}
}
@@ -740,8 +857,11 @@ impl ReplicationStats {
let mut interval = interval(Duration::from_secs(2));
loop {
interval.tick().await;
let current = get_global_replication_pool()
.map(|pool| pool.active_workers() + pool.active_lrg_workers() + pool.active_mrf_workers())
.unwrap_or(0);
let mut workers = workers_clone.lock().await;
workers.update();
workers.update(current);
}
});
@@ -925,14 +1045,26 @@ impl ReplicationStats {
/// Get replication metrics for all buckets
pub async fn get_all(&self) -> HashMap<String, BucketReplicationStats> {
let cache = self.cache.read().await;
let mut result = HashMap::new();
let mut result = HashMap::with_capacity(cache.len());
for (bucket, stats) in cache.iter() {
let mut cloned_stats = stats.clone_stats();
// Add queue statistics
result.insert(bucket.clone(), stats.clone_stats());
}
drop(cache);
{
let q_cache = self.q_cache.lock().await;
cloned_stats.q_stat = q_cache.get_bucket_stats(bucket);
result.insert(bucket.clone(), cloned_stats);
for (bucket, queue_stats) in &q_cache.bucket_stats {
let bucket_stats = result.entry(bucket.clone()).or_insert_with(BucketReplicationStats::new);
bucket_stats.q_stat = queue_stats.snapshot();
}
}
{
let p_cache = self.p_cache.lock().await;
for bucket in p_cache.bucket_stats.keys() {
result.entry(bucket.clone()).or_insert_with(BucketReplicationStats::new);
}
}
result
@@ -1114,12 +1246,12 @@ impl ReplicationStats {
let stats = q_cache
.bucket_stats
.entry(bucket.to_string())
.or_insert_with(InQueueStats::default);
stats.now_bytes.fetch_add(size, Ordering::Relaxed);
stats.now_count.fetch_add(1, Ordering::Relaxed);
.or_insert_with(InQueueMetric::default);
stats.curr.now_bytes.fetch_add(size, Ordering::Relaxed);
stats.curr.now_count.fetch_add(1, Ordering::Relaxed);
q_cache.sr_queue_stats.now_bytes.fetch_add(size, Ordering::Relaxed);
q_cache.sr_queue_stats.now_count.fetch_add(1, Ordering::Relaxed);
q_cache.sr_queue_stats.curr.now_bytes.fetch_add(size, Ordering::Relaxed);
q_cache.sr_queue_stats.curr.now_count.fetch_add(1, Ordering::Relaxed);
}
/// Decrease queue statistics
@@ -1128,12 +1260,12 @@ impl ReplicationStats {
let stats = q_cache
.bucket_stats
.entry(bucket.to_string())
.or_insert_with(InQueueStats::default);
stats.now_bytes.fetch_sub(size, Ordering::Relaxed);
stats.now_count.fetch_sub(1, Ordering::Relaxed);
.or_insert_with(InQueueMetric::default);
stats.curr.now_bytes.fetch_sub(size, Ordering::Relaxed);
stats.curr.now_count.fetch_sub(1, Ordering::Relaxed);
q_cache.sr_queue_stats.now_bytes.fetch_sub(size, Ordering::Relaxed);
q_cache.sr_queue_stats.now_count.fetch_sub(1, Ordering::Relaxed);
q_cache.sr_queue_stats.curr.now_bytes.fetch_sub(size, Ordering::Relaxed);
q_cache.sr_queue_stats.curr.now_count.fetch_sub(1, Ordering::Relaxed);
}
/// Increase proxy metrics
@@ -1166,6 +1298,51 @@ mod tests {
assert_eq!(workers.curr, 0);
}
#[test]
fn test_in_queue_metric_observe_updates_rolling_stats() {
let mut metric = InQueueMetric::default();
metric.curr.now_bytes.store(128, Ordering::Relaxed);
metric.curr.now_count.store(4, Ordering::Relaxed);
metric.observe(Instant::now());
metric.curr.now_bytes.store(256, Ordering::Relaxed);
metric.curr.now_count.store(6, Ordering::Relaxed);
metric.observe(Instant::now());
assert_eq!(metric.curr.bytes, 256);
assert_eq!(metric.curr.count, 6);
assert_eq!(metric.max.bytes, 256);
assert_eq!(metric.max.count, 6);
assert_eq!(metric.last_minute.bytes, 192);
assert_eq!(metric.last_minute.count, 5);
}
#[test]
fn test_fail_stats_recent_since_tracks_windows() {
let mut stats = FailStats::default();
stats.add_size(64, None);
stats.add_size(32, None);
let last_minute = stats.recent_since(Duration::from_secs(60));
let last_hour = stats.recent_since(Duration::from_secs(60 * 60));
assert_eq!(last_minute.count, 2);
assert_eq!(last_minute.size, 96);
assert_eq!(last_hour.count, 2);
assert_eq!(last_hour.size, 96);
}
#[test]
fn test_active_worker_stat_update_tracks_rolling_avg_and_max() {
let mut stats = ActiveWorkerStat::default();
stats.update(2);
stats.update(6);
stats.update(4);
assert_eq!(stats.curr, 4);
assert_eq!(stats.max, 6);
assert_eq!(stats.avg, 4.0);
}
#[tokio::test]
async fn test_delete_bucket_stats() {
let stats = ReplicationStats::new();
@@ -1218,6 +1395,15 @@ mod tests {
assert_eq!(stat.replicated_count, 1);
}
#[tokio::test]
async fn test_get_all_includes_proxy_only_bucket() {
let stats = ReplicationStats::new();
stats.inc_proxy("proxy-only-bucket", "HeadObject", false).await;
let all = stats.get_all().await;
assert!(all.contains_key("proxy-only-bucket"));
}
#[test]
fn test_sr_stats() {
let sr_stats = SRStats::new();