Skip to main content

rc_core/admin/
replication.rs

1//! Typed contracts for read-only replication inspection.
2
3use std::collections::BTreeMap;
4
5use async_trait::async_trait;
6use jiff::Timestamp;
7use serde::{Deserialize, Serialize};
8use serde_json::Value;
9
10use crate::Result;
11
12/// Maximum encoded size accepted for one replication diff response.
13pub const MAX_REPLICATION_DIFF_RESPONSE_BYTES: usize = 8 * 1024 * 1024;
14/// Maximum encoded size accepted for metrics and MRF responses.
15pub const MAX_REPLICATION_INSPECTION_RESPONSE_BYTES: usize = 8 * 1024 * 1024;
16
17/// Scope of a replication observation. Unknown values are retained for forward compatibility.
18#[derive(Debug, Clone, PartialEq, Eq)]
19pub enum ReplicationMetricScope {
20    Unavailable,
21    NodeLocal,
22    ClusterAggregated,
23    PartialCluster,
24    Unknown(String),
25}
26
27impl ReplicationMetricScope {
28    pub fn as_str(&self) -> &str {
29        match self {
30            Self::Unavailable => "unavailable",
31            Self::NodeLocal => "node_local",
32            Self::ClusterAggregated => "cluster_aggregated",
33            Self::PartialCluster => "partial_cluster",
34            Self::Unknown(value) => value,
35        }
36    }
37}
38
39impl Serialize for ReplicationMetricScope {
40    fn serialize<S: serde::Serializer>(
41        &self,
42        serializer: S,
43    ) -> std::result::Result<S::Ok, S::Error> {
44        serializer.serialize_str(self.as_str())
45    }
46}
47
48impl<'de> Deserialize<'de> for ReplicationMetricScope {
49    fn deserialize<D: serde::Deserializer<'de>>(
50        deserializer: D,
51    ) -> std::result::Result<Self, D::Error> {
52        let value = String::deserialize(deserializer)?;
53        Ok(match value.as_str() {
54            "unavailable" => Self::Unavailable,
55            "node_local" => Self::NodeLocal,
56            "cluster_aggregated" => Self::ClusterAggregated,
57            "partial_cluster" => Self::PartialCluster,
58            _ => Self::Unknown(value),
59        })
60    }
61}
62
63#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
64pub struct ReplicationCountSize {
65    pub count: u64,
66    #[serde(rename = "bytes", alias = "size")]
67    pub size: u64,
68    #[serde(flatten, default)]
69    pub extra: BTreeMap<String, Value>,
70}
71
72#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
73pub struct ReplicationQueueMetric {
74    pub curr: ReplicationCountSize,
75    pub avg: ReplicationCountSize,
76    pub max: ReplicationCountSize,
77    pub last_minute: ReplicationCountSize,
78    #[serde(flatten, default)]
79    pub extra: BTreeMap<String, Value>,
80}
81
82#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
83pub struct ReplicationLatencyMetric {
84    pub avg: f64,
85    pub curr: f64,
86    pub max: f64,
87    #[serde(flatten, default)]
88    pub extra: BTreeMap<String, Value>,
89}
90
91#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
92pub struct ReplicationTransferRate {
93    pub avg: f64,
94    pub curr: f64,
95    pub peak: f64,
96    #[serde(flatten, default)]
97    pub extra: BTreeMap<String, Value>,
98}
99
100#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
101pub struct ReplicationTargetMetric {
102    pub replicated_size: u64,
103    pub replicated_count: u64,
104    pub failed: ReplicationCountSize,
105    #[serde(default)]
106    pub fail_stats: Option<ReplicationCountSize>,
107    pub latency: ReplicationLatencyMetric,
108    pub xfer_rate_lrg: ReplicationTransferRate,
109    pub xfer_rate_sml: ReplicationTransferRate,
110    pub bandwidth_limit_bytes_per_sec: u64,
111    pub current_bandwidth_bytes_per_sec: f64,
112    #[serde(default)]
113    pub latency_scope: Option<ReplicationMetricScope>,
114    #[serde(default)]
115    pub bandwidth_scope: Option<ReplicationMetricScope>,
116    #[serde(flatten, default)]
117    pub extra: BTreeMap<String, Value>,
118}
119
120#[derive(Debug, Clone, PartialEq, Serialize)]
121pub struct ReplicationMetrics {
122    pub stats: BTreeMap<String, ReplicationTargetMetric>,
123    pub replica_size: u64,
124    pub replica_count: u64,
125    pub replicated_size: u64,
126    pub replicated_count: u64,
127    pub q_stat: ReplicationQueueMetric,
128    #[serde(default)]
129    pub provider_available: Option<bool>,
130    #[serde(default)]
131    pub cluster_complete: Option<bool>,
132    #[serde(default)]
133    pub observed_node_count: Option<u32>,
134    #[serde(default)]
135    pub expected_node_count: Option<u32>,
136    #[serde(default)]
137    pub queue_scope: Option<ReplicationMetricScope>,
138    #[serde(flatten, default)]
139    pub extra: BTreeMap<String, Value>,
140}
141
142#[derive(Debug, Default)]
143enum WireField<T> {
144    #[default]
145    Missing,
146    Present(T),
147}
148
149impl<T> WireField<T> {
150    fn is_present(&self) -> bool {
151        matches!(self, Self::Present(_))
152    }
153
154    fn require(self, name: &str) -> std::result::Result<T, String> {
155        match self {
156            Self::Present(value) => Ok(value),
157            Self::Missing => Err(format!("missing field `{name}`")),
158        }
159    }
160}
161
162impl<'de, T> Deserialize<'de> for WireField<T>
163where
164    T: Deserialize<'de>,
165{
166    fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
167    where
168        D: serde::Deserializer<'de>,
169    {
170        T::deserialize(deserializer).map(Self::Present)
171    }
172}
173
174#[derive(Debug, Clone, Copy)]
175struct WireCounter(u64);
176
177impl<'de> Deserialize<'de> for WireCounter {
178    fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
179    where
180        D: serde::Deserializer<'de>,
181    {
182        // Keep the original token so decimal counters cannot be silently rounded
183        // through f64 before their integer semantics are validated.
184        let raw = <&serde_json::value::RawValue>::deserialize(deserializer)?;
185        parse_wire_counter(raw.get())
186            .map(WireCounter)
187            .map_err(serde::de::Error::custom)
188    }
189}
190
191fn parse_wire_counter(raw: &str) -> std::result::Result<u64, &'static str> {
192    let bytes = raw.as_bytes();
193    if bytes.is_empty() || bytes[0] == b'-' {
194        return Err("counter cannot be negative");
195    }
196    if !bytes[0].is_ascii_digit() {
197        return Err("counter must be a JSON number");
198    }
199
200    let exponent_start = bytes
201        .iter()
202        .position(|byte| matches!(byte, b'e' | b'E'))
203        .unwrap_or(bytes.len());
204    let mantissa = &bytes[..exponent_start];
205    let decimal = mantissa.iter().position(|byte| *byte == b'.');
206    let (integer_digits, fraction_digits) = match decimal {
207        Some(index) => (&mantissa[..index], &mantissa[index + 1..]),
208        None => (mantissa, &[][..]),
209    };
210    if integer_digits.is_empty()
211        || !integer_digits.iter().all(u8::is_ascii_digit)
212        || !fraction_digits.iter().all(u8::is_ascii_digit)
213    {
214        return Err("counter must be a JSON number");
215    }
216
217    let digits = integer_digits.iter().chain(fraction_digits);
218    if digits.clone().all(|digit| *digit == b'0') {
219        return Ok(0);
220    }
221
222    let exponent = if exponent_start == bytes.len() {
223        0_i64
224    } else {
225        parse_counter_exponent(&bytes[exponent_start + 1..])?
226    };
227    let fraction_len = i64::try_from(fraction_digits.len())
228        .map_err(|_| "counter has too many fractional digits")?;
229    let scale = exponent.saturating_sub(fraction_len);
230    let total_digits = integer_digits.len() + fraction_digits.len();
231    let (kept_digits, appended_zeros) = if scale < 0 {
232        let removed_digits = usize::try_from(scale.unsigned_abs())
233            .map_err(|_| "counter contains a fractional value")?;
234        if removed_digits > total_digits {
235            return Err("counter contains a fractional value");
236        }
237        let kept_digits = total_digits - removed_digits;
238        if integer_digits
239            .iter()
240            .chain(fraction_digits)
241            .skip(kept_digits)
242            .any(|digit| *digit != b'0')
243        {
244            return Err("counter contains a fractional value");
245        }
246        (kept_digits, 0_usize)
247    } else {
248        let appended_zeros = usize::try_from(scale).map_err(|_| "counter exceeds the u64 range")?;
249        if appended_zeros > 20 {
250            return Err("counter exceeds the u64 range");
251        }
252        (total_digits, appended_zeros)
253    };
254
255    let mut counter = 0_u64;
256    for digit in integer_digits
257        .iter()
258        .chain(fraction_digits)
259        .take(kept_digits)
260    {
261        counter = counter
262            .checked_mul(10)
263            .and_then(|value| value.checked_add(u64::from(*digit - b'0')))
264            .ok_or("counter exceeds the u64 range")?;
265    }
266    for _ in 0..appended_zeros {
267        counter = counter
268            .checked_mul(10)
269            .ok_or("counter exceeds the u64 range")?;
270    }
271    Ok(counter)
272}
273
274fn parse_counter_exponent(raw: &[u8]) -> std::result::Result<i64, &'static str> {
275    let (negative, digits) = match raw.first() {
276        Some(b'+') => (false, &raw[1..]),
277        Some(b'-') => (true, &raw[1..]),
278        Some(_) => (false, raw),
279        None => return Err("counter exponent is missing"),
280    };
281    if digits.is_empty() || !digits.iter().all(u8::is_ascii_digit) {
282        return Err("counter exponent is invalid");
283    }
284    let exponent = digits.iter().fold(0_i64, |value, digit| {
285        value
286            .saturating_mul(10)
287            .saturating_add(i64::from(*digit - b'0'))
288    });
289    Ok(if negative {
290        exponent.saturating_neg()
291    } else {
292        exponent
293    })
294}
295
296#[derive(Debug, Deserialize)]
297struct MinioCountSizeWire {
298    count: WireCounter,
299    #[serde(rename = "bytes")]
300    size: WireCounter,
301}
302
303impl From<MinioCountSizeWire> for ReplicationCountSize {
304    fn from(value: MinioCountSizeWire) -> Self {
305        Self {
306            count: value.count.0,
307            size: value.size.0,
308            extra: BTreeMap::new(),
309        }
310    }
311}
312
313#[derive(Debug, Deserialize)]
314struct MinioTimedErrStatsWire {
315    #[serde(rename = "lastMinute")]
316    last_minute: MinioCountSizeWire,
317    #[serde(rename = "lastHour")]
318    last_hour: MinioCountSizeWire,
319    totals: MinioCountSizeWire,
320}
321
322impl MinioTimedErrStatsWire {
323    fn into_totals(self) -> ReplicationCountSize {
324        let Self {
325            last_minute,
326            last_hour,
327            totals,
328        } = self;
329        let _ = (last_minute, last_hour);
330        totals.into()
331    }
332}
333
334#[derive(Debug, Deserialize)]
335struct MinioQueueMetricWire {
336    curr: MinioCountSizeWire,
337    avg: MinioCountSizeWire,
338    #[serde(default)]
339    max: WireField<MinioCountSizeWire>,
340    #[serde(default)]
341    peak: WireField<MinioCountSizeWire>,
342}
343
344impl MinioQueueMetricWire {
345    fn try_into_metric(self) -> std::result::Result<ReplicationQueueMetric, String> {
346        let maximum = match (self.max, self.peak) {
347            (WireField::Present(max), WireField::Present(peak)) => {
348                let max: ReplicationCountSize = max.into();
349                let peak: ReplicationCountSize = peak.into();
350                if max != peak {
351                    return Err("fields `queued.max` and `queued.peak` conflict".into());
352                }
353                max
354            }
355            (WireField::Present(max), WireField::Missing) => max.into(),
356            (WireField::Missing, WireField::Present(peak)) => peak.into(),
357            (WireField::Missing, WireField::Missing) => {
358                return Err(
359                    "missing field `queued.peak` or compatibility field `queued.max`".into(),
360                );
361            }
362        };
363
364        Ok(ReplicationQueueMetric {
365            curr: self.curr.into(),
366            avg: self.avg.into(),
367            max: maximum,
368            last_minute: ReplicationCountSize::default(),
369            extra: BTreeMap::new(),
370        })
371    }
372}
373
374#[derive(Debug, Deserialize)]
375struct MinioTargetMetricWire {
376    #[serde(rename = "replicationCount", default)]
377    replicated_count: WireField<WireCounter>,
378    #[serde(rename = "completedReplicationSize", default)]
379    replicated_size: WireField<WireCounter>,
380    #[serde(rename = "limitInBits", default)]
381    bandwidth_limit: WireField<WireCounter>,
382    #[serde(rename = "currentBandwidth", default)]
383    current_bandwidth: WireField<f64>,
384    #[serde(default)]
385    failed: WireField<MinioTimedErrStatsWire>,
386    #[serde(rename = "pendingReplicationSize", default)]
387    pending_size: WireField<WireCounter>,
388    #[serde(rename = "replicaSize", default)]
389    replica_size: WireField<WireCounter>,
390    #[serde(rename = "failedReplicationSize", default)]
391    failed_size: WireField<WireCounter>,
392    #[serde(rename = "pendingReplicationCount", default)]
393    pending_count: WireField<WireCounter>,
394    #[serde(rename = "failedReplicationCount", default)]
395    failed_count: WireField<WireCounter>,
396    #[serde(flatten, default)]
397    extra: BTreeMap<String, Value>,
398}
399
400impl MinioTargetMetricWire {
401    fn try_into_metric(self) -> std::result::Result<ReplicationTargetMetric, String> {
402        let Self {
403            replicated_count,
404            replicated_size,
405            bandwidth_limit,
406            current_bandwidth,
407            failed,
408            pending_size,
409            replica_size,
410            failed_size,
411            pending_count,
412            failed_count,
413            extra,
414        } = self;
415        let failed = minio_failed_totals(failed);
416        validate_redundant_failed(&failed, &failed_count, &failed_size)?;
417        let current_bandwidth = match current_bandwidth {
418            WireField::Missing => 0.0,
419            WireField::Present(value) if value.is_finite() && value >= 0.0 => value,
420            WireField::Present(_) => {
421                return Err("field `currentBandwidth` must be finite and non-negative".into());
422            }
423        };
424        let _ = (pending_size, replica_size, pending_count);
425
426        Ok(ReplicationTargetMetric {
427            replicated_size: wire_counter_or_zero(replicated_size),
428            replicated_count: wire_counter_or_zero(replicated_count),
429            failed,
430            fail_stats: None,
431            latency: ReplicationLatencyMetric::default(),
432            xfer_rate_lrg: ReplicationTransferRate::default(),
433            xfer_rate_sml: ReplicationTransferRate::default(),
434            bandwidth_limit_bytes_per_sec: wire_counter_or_zero(bandwidth_limit),
435            current_bandwidth_bytes_per_sec: current_bandwidth,
436            latency_scope: Some(ReplicationMetricScope::Unavailable),
437            bandwidth_scope: None,
438            extra,
439        })
440    }
441}
442
443#[derive(Debug, Deserialize)]
444struct ReplicationMetricsWireEnvelope {
445    #[serde(rename = "Stats", default)]
446    minio_stats: WireField<Option<BTreeMap<String, MinioTargetMetricWire>>>,
447    #[serde(rename = "completedReplicationSize", default)]
448    minio_replicated_size: WireField<WireCounter>,
449    #[serde(rename = "replicaSize", default)]
450    minio_replica_size: WireField<WireCounter>,
451    #[serde(rename = "replicaCount", default)]
452    minio_replica_count: WireField<WireCounter>,
453    #[serde(rename = "replicationCount", default)]
454    minio_replicated_count: WireField<WireCounter>,
455    #[serde(rename = "failed", default)]
456    minio_failed: WireField<MinioTimedErrStatsWire>,
457    #[serde(rename = "queued", default)]
458    minio_queue: WireField<MinioQueueMetricWire>,
459    #[serde(rename = "pendingReplicationSize", default)]
460    minio_pending_size: WireField<WireCounter>,
461    #[serde(rename = "failedReplicationSize", default)]
462    minio_failed_size: WireField<WireCounter>,
463    #[serde(rename = "pendingReplicationCount", default)]
464    minio_pending_count: WireField<WireCounter>,
465    #[serde(rename = "failedReplicationCount", default)]
466    minio_failed_count: WireField<WireCounter>,
467    #[serde(rename = "stats", default)]
468    legacy_stats: WireField<BTreeMap<String, ReplicationTargetMetric>>,
469    #[serde(rename = "replica_size", default)]
470    legacy_replica_size: WireField<u64>,
471    #[serde(rename = "replica_count", default)]
472    legacy_replica_count: WireField<u64>,
473    #[serde(rename = "replicated_size", default)]
474    legacy_replicated_size: WireField<u64>,
475    #[serde(rename = "replicated_count", default)]
476    legacy_replicated_count: WireField<u64>,
477    #[serde(rename = "q_stat", default)]
478    legacy_queue: WireField<ReplicationQueueMetric>,
479    #[serde(default)]
480    provider_available: Option<bool>,
481    #[serde(default)]
482    cluster_complete: Option<bool>,
483    #[serde(default)]
484    observed_node_count: Option<u32>,
485    #[serde(default)]
486    expected_node_count: Option<u32>,
487    #[serde(default)]
488    queue_scope: Option<ReplicationMetricScope>,
489    #[serde(flatten, default)]
490    extra: BTreeMap<String, Value>,
491}
492
493impl ReplicationMetricsWireEnvelope {
494    fn has_minio_fields(&self) -> bool {
495        self.minio_stats.is_present()
496            || self.minio_replicated_size.is_present()
497            || self.minio_replica_size.is_present()
498            || self.minio_replica_count.is_present()
499            || self.minio_replicated_count.is_present()
500            || self.minio_failed.is_present()
501            || self.minio_queue.is_present()
502            || self.minio_pending_size.is_present()
503            || self.minio_failed_size.is_present()
504            || self.minio_pending_count.is_present()
505            || self.minio_failed_count.is_present()
506    }
507
508    fn has_legacy_fields(&self) -> bool {
509        self.legacy_stats.is_present()
510            || self.legacy_replica_size.is_present()
511            || self.legacy_replica_count.is_present()
512            || self.legacy_replicated_size.is_present()
513            || self.legacy_replicated_count.is_present()
514            || self.legacy_queue.is_present()
515    }
516
517    fn try_into_metrics(self) -> std::result::Result<ReplicationMetrics, String> {
518        match (
519            self.minio_stats.is_present(),
520            self.legacy_stats.is_present(),
521        ) {
522            (true, true) => Err("replication metrics response mixes `Stats` and `stats`".into()),
523            (false, false) => {
524                Err("missing replication metrics discriminator `Stats` or `stats`".into())
525            }
526            (true, false) => {
527                if self.has_legacy_fields() {
528                    return Err("MinIO replication metrics contain legacy fields".into());
529                }
530                self.try_into_minio_metrics()
531            }
532            (false, true) => {
533                if self.has_minio_fields() {
534                    return Err("legacy replication metrics contain MinIO fields".into());
535                }
536                self.try_into_legacy_metrics()
537            }
538        }
539    }
540
541    fn try_into_minio_metrics(self) -> std::result::Result<ReplicationMetrics, String> {
542        let stats = self
543            .minio_stats
544            .require("Stats")?
545            .unwrap_or_default()
546            .into_iter()
547            .map(|(arn, target)| target.try_into_metric().map(|target| (arn, target)))
548            .collect::<std::result::Result<BTreeMap<_, _>, _>>()?;
549        let failed = minio_failed_totals(self.minio_failed);
550        validate_redundant_failed(&failed, &self.minio_failed_count, &self.minio_failed_size)?;
551        let queue = self.minio_queue.require("queued")?.try_into_metric()?;
552        let _ = (self.minio_pending_size, self.minio_pending_count, failed);
553
554        Ok(ReplicationMetrics {
555            stats,
556            replica_size: wire_counter_or_zero(self.minio_replica_size),
557            replica_count: wire_counter_or_zero(self.minio_replica_count),
558            replicated_size: wire_counter_or_zero(self.minio_replicated_size),
559            replicated_count: wire_counter_or_zero(self.minio_replicated_count),
560            q_stat: queue,
561            provider_available: self.provider_available,
562            cluster_complete: self.cluster_complete,
563            observed_node_count: self.observed_node_count,
564            expected_node_count: self.expected_node_count,
565            queue_scope: self.queue_scope,
566            extra: self.extra,
567        })
568    }
569
570    fn try_into_legacy_metrics(self) -> std::result::Result<ReplicationMetrics, String> {
571        Ok(ReplicationMetrics {
572            stats: self.legacy_stats.require("stats")?,
573            replica_size: self.legacy_replica_size.require("replica_size")?,
574            replica_count: self.legacy_replica_count.require("replica_count")?,
575            replicated_size: self.legacy_replicated_size.require("replicated_size")?,
576            replicated_count: self.legacy_replicated_count.require("replicated_count")?,
577            q_stat: self.legacy_queue.require("q_stat")?,
578            provider_available: self.provider_available,
579            cluster_complete: self.cluster_complete,
580            observed_node_count: self.observed_node_count,
581            expected_node_count: self.expected_node_count,
582            queue_scope: self.queue_scope,
583            extra: self.extra,
584        })
585    }
586}
587
588fn wire_counter_or_zero(field: WireField<WireCounter>) -> u64 {
589    match field {
590        WireField::Missing => 0,
591        WireField::Present(value) => value.0,
592    }
593}
594
595fn minio_failed_totals(field: WireField<MinioTimedErrStatsWire>) -> ReplicationCountSize {
596    match field {
597        WireField::Missing => ReplicationCountSize::default(),
598        WireField::Present(value) => value.into_totals(),
599    }
600}
601
602fn validate_redundant_failed(
603    totals: &ReplicationCountSize,
604    count: &WireField<WireCounter>,
605    size: &WireField<WireCounter>,
606) -> std::result::Result<(), String> {
607    if matches!(count, WireField::Present(value) if value.0 != totals.count)
608        || matches!(size, WireField::Present(value) if value.0 != totals.size)
609    {
610        return Err("replication metrics contain inconsistent failed totals".into());
611    }
612    Ok(())
613}
614
615impl<'de> Deserialize<'de> for ReplicationMetrics {
616    fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
617    where
618        D: serde::Deserializer<'de>,
619    {
620        ReplicationMetricsWireEnvelope::deserialize(deserializer)?
621            .try_into_metrics()
622            .map_err(serde::de::Error::custom)
623    }
624}
625
626#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
627pub struct ReplicationMrfTarget {
628    #[serde(rename = "ARN")]
629    pub arn: String,
630    #[serde(rename = "FailedCount")]
631    pub failed_count: u64,
632    #[serde(rename = "FailedSize")]
633    pub failed_size: u64,
634    #[serde(rename = "ObservationScope")]
635    pub observation_scope: ReplicationMetricScope,
636    #[serde(flatten, default)]
637    pub extra: BTreeMap<String, Value>,
638}
639
640#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
641pub struct ReplicationMrf {
642    #[serde(rename = "Bucket")]
643    pub bucket: String,
644    #[serde(rename = "Targets")]
645    pub targets: Vec<ReplicationMrfTarget>,
646    #[serde(rename = "TotalFailedCount")]
647    pub total_failed_count: u64,
648    #[serde(rename = "TotalFailedSize")]
649    pub total_failed_size: u64,
650    #[serde(rename = "QueuedCount")]
651    pub queued_count: u64,
652    #[serde(rename = "QueuedSize")]
653    pub queued_size: u64,
654    #[serde(rename = "PerObjectEntriesAvailable")]
655    pub per_object_entries_available: bool,
656    #[serde(rename = "RuntimeStatsAvailable")]
657    pub runtime_stats_available: bool,
658    #[serde(rename = "ClusterComplete")]
659    pub cluster_complete: bool,
660    #[serde(rename = "ObservedNodeCount")]
661    pub observed_node_count: u32,
662    #[serde(rename = "ExpectedNodeCount")]
663    pub expected_node_count: u32,
664    #[serde(rename = "DurableBacklogAvailable")]
665    pub durable_backlog_available: bool,
666    #[serde(rename = "DurableCount")]
667    pub durable_count: u64,
668    #[serde(rename = "DurableSize")]
669    pub durable_size: u64,
670    #[serde(rename = "PerTargetDurableEntriesAvailable")]
671    pub per_target_durable_entries_available: bool,
672    #[serde(flatten, default)]
673    pub extra: BTreeMap<String, Value>,
674}
675
676#[async_trait]
677pub trait ReplicationInspectionApi: Send + Sync {
678    async fn replication_metrics(&self, bucket: &str) -> Result<ReplicationMetrics>;
679    async fn replication_mrf(&self, bucket: &str) -> Result<ReplicationMrf>;
680}
681
682/// A bounded, on-demand scan of object versions that have not replicated.
683#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
684pub struct ReplicationDiff {
685    #[serde(rename = "Entries")]
686    pub entries: Vec<ReplicationDiffEntry>,
687    #[serde(rename = "IsTruncated")]
688    pub is_truncated: bool,
689    #[serde(rename = "ScannedVersions")]
690    pub scanned_versions: usize,
691    #[serde(flatten, default)]
692    pub extra: BTreeMap<String, Value>,
693}
694
695/// One pending or failed object version returned by a replication diff.
696#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
697pub struct ReplicationDiffEntry {
698    #[serde(rename = "Object")]
699    pub object: String,
700    #[serde(rename = "VersionID")]
701    pub version_id: Option<String>,
702    #[serde(rename = "Size")]
703    pub size_bytes: u64,
704    #[serde(rename = "IsDeleteMarker")]
705    pub delete_marker: bool,
706    #[serde(rename = "ReplicationStatus")]
707    pub replication_status: String,
708    #[serde(rename = "LastModified")]
709    pub last_modified: Option<Timestamp>,
710    #[serde(flatten, default)]
711    pub extra: BTreeMap<String, Value>,
712}
713
714/// Read-only RustFS replication diff operations.
715#[async_trait]
716pub trait ReplicationDiffApi: Send + Sync {
717    /// Scan a bucket, optionally below a prefix, for pending or failed versions.
718    async fn replication_diff(&self, bucket: &str, prefix: Option<&str>)
719    -> Result<ReplicationDiff>;
720}
721
722#[cfg(test)]
723mod tests {
724    use super::*;
725
726    #[test]
727    fn response_preserves_unknown_fields_and_typed_entries() {
728        let response: ReplicationDiff = serde_json::from_str(
729            r#"{
730                "Entries": [{
731                    "Object": "reports/a.json",
732                    "VersionID": "v1",
733                    "Size": 42,
734                    "IsDeleteMarker": false,
735                    "ReplicationStatus": "FAILED",
736                    "LastModified": "2026-07-21T04:00:00Z",
737                    "TargetDetail": {"attempts": 2}
738                }],
739                "IsTruncated": true,
740                "ScannedVersions": 10000,
741                "ServerRevision": 7
742            }"#,
743        )
744        .expect("typed replication diff");
745
746        assert_eq!(response.entries[0].size_bytes, 42);
747        assert_eq!(response.entries[0].version_id.as_deref(), Some("v1"));
748        assert_eq!(response.entries[0].extra["TargetDetail"]["attempts"], 2);
749        assert_eq!(response.extra["ServerRevision"], 7);
750    }
751
752    #[test]
753    fn response_accepts_delete_marker_without_version_or_timestamp() {
754        let response: ReplicationDiff = serde_json::from_str(
755            r#"{
756                "Entries": [{
757                    "Object": "removed.txt",
758                    "VersionID": null,
759                    "Size": 0,
760                    "IsDeleteMarker": true,
761                    "ReplicationStatus": "PENDING",
762                    "LastModified": null
763                }],
764                "IsTruncated": false,
765                "ScannedVersions": 1
766            }"#,
767        )
768        .expect("delete marker diff");
769
770        assert!(response.entries[0].delete_marker);
771        assert!(response.entries[0].version_id.is_none());
772        assert!(response.entries[0].last_modified.is_none());
773    }
774
775    #[test]
776    fn response_rejects_negative_sizes_and_malformed_timestamps() {
777        for payload in [
778            r#"{"Entries":[{"Object":"a","VersionID":null,"Size":-1,"IsDeleteMarker":false,"ReplicationStatus":"FAILED","LastModified":null}],"IsTruncated":false,"ScannedVersions":1}"#,
779            r#"{"Entries":[{"Object":"a","VersionID":null,"Size":1,"IsDeleteMarker":false,"ReplicationStatus":"FAILED","LastModified":"yesterday"}],"IsTruncated":false,"ScannedVersions":1}"#,
780        ] {
781            assert!(serde_json::from_str::<ReplicationDiff>(payload).is_err());
782        }
783    }
784
785    #[test]
786    fn response_requires_scan_completeness_fields() {
787        let payload = r#"{"Entries":[]}"#;
788        assert!(serde_json::from_str::<ReplicationDiff>(payload).is_err());
789    }
790
791    #[test]
792    fn metrics_distinguish_legacy_metadata_and_preserve_unknown_scope() {
793        let legacy: ReplicationMetrics = serde_json::from_str(r#"{"stats":{},"replica_size":0,"replica_count":0,"replicated_size":0,"replicated_count":0,"q_stat":{"curr":{"count":0,"size":0},"avg":{"count":0,"size":0},"max":{"count":0,"size":0},"last_minute":{"count":0,"size":0}}}"#).expect("legacy metrics");
794        assert_eq!(legacy.provider_available, None);
795        assert_eq!(legacy.cluster_complete, None);
796
797        let current: ReplicationMetrics = serde_json::from_str(r#"{"stats":{},"replica_size":0,"replica_count":0,"replicated_size":0,"replicated_count":0,"q_stat":{"curr":{"count":0,"size":0},"avg":{"count":0,"size":0},"max":{"count":0,"size":0},"last_minute":{"count":0,"size":0}},"provider_available":true,"cluster_complete":false,"observed_node_count":1,"expected_node_count":2,"queue_scope":"future_scope"}"#).expect("current metrics");
798        assert_eq!(current.provider_available, Some(true));
799        assert_eq!(
800            current.queue_scope,
801            Some(ReplicationMetricScope::Unknown("future_scope".into()))
802        );
803    }
804
805    #[test]
806    fn metrics_and_mrf_reject_negative_counters() {
807        let metrics = r#"{"stats":{},"replica_size":-1,"replica_count":0,"replicated_size":0,"replicated_count":0,"q_stat":{"curr":{"count":0,"size":0},"avg":{"count":0,"size":0},"max":{"count":0,"size":0},"last_minute":{"count":0,"size":0}}}"#;
808        assert!(serde_json::from_str::<ReplicationMetrics>(metrics).is_err());
809        let mrf = r#"{"Bucket":"b","Targets":[],"TotalFailedCount":-1,"TotalFailedSize":0,"QueuedCount":0,"QueuedSize":0,"PerObjectEntriesAvailable":false,"RuntimeStatsAvailable":true,"ClusterComplete":false,"ObservedNodeCount":1,"ExpectedNodeCount":2,"DurableBacklogAvailable":false,"DurableCount":0,"DurableSize":0,"PerTargetDurableEntriesAvailable":false}"#;
810        assert!(serde_json::from_str::<ReplicationMrf>(mrf).is_err());
811    }
812
813    #[test]
814    fn metrics_decode_captured_minio_v1_wire_response() {
815        let metrics: ReplicationMetrics = serde_json::from_str(include_str!(
816            "../../tests/fixtures/replication_metrics_minio_v1.json"
817        ))
818        .expect("captured MinIO-compatible metrics");
819
820        let target = metrics
821            .stats
822            .get("arn:minio:replication:us-east-1:00000000-0000-0000-0000-000000000000:destination")
823            .expect("captured target");
824        assert_eq!(metrics.replicated_count, 1);
825        assert_eq!(metrics.replicated_size, 20);
826        assert_eq!(target.replicated_count, 1);
827        assert_eq!(target.replicated_size, 20);
828        assert_eq!(
829            target.latency_scope,
830            Some(ReplicationMetricScope::Unavailable)
831        );
832        assert_eq!(metrics.provider_available, Some(true));
833        assert_eq!(metrics.cluster_complete, Some(true));
834        assert_eq!(metrics.observed_node_count, Some(1));
835        assert_eq!(metrics.expected_node_count, Some(1));
836    }
837
838    #[test]
839    fn metrics_use_timed_totals_and_preserve_minio_extensions() {
840        let metrics: ReplicationMetrics = serde_json::from_str(
841            r#"{
842                "Stats":{"arn:target":{
843                    "replicationCount":2,"completedReplicationSize":30,
844                    "limitInBits":8000,"currentBandwidth":125.5,
845                    "failed":{
846                        "lastMinute":{"count":1.0,"bytes":10.0},
847                        "lastHour":{"count":4.0,"bytes":40.0},
848                        "totals":{"count":9.0,"bytes":90.0}},
849                    "failedReplicationCount":9,"failedReplicationSize":90,
850                    "TargetFuture":{"token":"value"}}},
851                "completedReplicationSize":30,"replicaSize":4,
852                "replicaCount":3,"replicationCount":2,
853                "failed":{
854                    "lastMinute":{"count":1.0,"bytes":10.0},
855                    "lastHour":{"count":4.0,"bytes":40.0},
856                    "totals":{"count":9.0,"bytes":90.0}},
857                "queued":{
858                    "curr":{"count":1.0,"bytes":10.0},
859                    "avg":{"count":2.0,"bytes":20.0},
860                    "peak":{"count":3.0,"bytes":30.0}},
861                "TopFuture":{"revision":7}}
862            "#,
863        )
864        .expect("MinIO metrics");
865
866        let target = &metrics.stats["arn:target"];
867        assert_eq!(target.failed.count, 9);
868        assert_eq!(target.failed.size, 90);
869        assert_eq!(target.bandwidth_limit_bytes_per_sec, 8000);
870        assert_eq!(target.current_bandwidth_bytes_per_sec, 125.5);
871        assert_eq!(target.extra["TargetFuture"]["token"], "value");
872        assert_eq!(metrics.q_stat.max.count, 3);
873        assert_eq!(metrics.extra["TopFuture"]["revision"], 7);
874    }
875
876    #[test]
877    fn metrics_accept_queue_peak_or_max_but_reject_conflicts() {
878        fn payload(queue_tail: &str) -> String {
879            format!(
880                r#"{{"Stats":null,"queued":{{
881                    "curr":{{"count":0,"bytes":0}},
882                    "avg":{{"count":0,"bytes":0}},
883                    {queue_tail}}}}}"#
884            )
885        }
886
887        for tail in [
888            r#""peak":{"count":2,"bytes":20}"#,
889            r#""max":{"count":2,"bytes":20}"#,
890            r#""max":{"count":2,"bytes":20},"peak":{"count":2,"bytes":20}"#,
891        ] {
892            let metrics: ReplicationMetrics =
893                serde_json::from_str(&payload(tail)).expect("compatible queue peak");
894            assert_eq!(metrics.q_stat.max.count, 2);
895            assert_eq!(metrics.q_stat.max.size, 20);
896        }
897
898        assert!(
899            serde_json::from_str::<ReplicationMetrics>(&payload(
900                r#""max":{"count":2,"bytes":20},"peak":{"count":3,"bytes":20}"#,
901            ))
902            .is_err()
903        );
904    }
905
906    #[test]
907    fn metrics_accept_omitempty_and_lossless_numeric_encodings() {
908        for (encoded, expected) in [
909            ("1e3", 1_000_u64),
910            ("1.5e1", 15),
911            ("1.2300e2", 123),
912            ("0e-400", 0),
913            ("1000000000000000.0", 1_000_000_000_000_000),
914            ("9007199254740991.0", 9_007_199_254_740_991),
915            ("9007199254740992.0", 9_007_199_254_740_992),
916            ("9223372036854775808.0", 9_223_372_036_854_775_808),
917            ("1e19", 10_000_000_000_000_000_000),
918            ("18446744073709551615", u64::MAX),
919            ("18446744073709551615.0", u64::MAX),
920        ] {
921            let payload = format!(
922                r#"{{"Stats":{{"arn:target":{{"replicationCount":1e3}}}},
923                    "queued":{{"curr":{{"count":0.0,"bytes":0.0}},
924                    "avg":{{"count":0,"bytes":0}},
925                    "peak":{{"count":{encoded},"bytes":42.0}}}}}}"#
926            );
927            let metrics: ReplicationMetrics =
928                serde_json::from_str(&payload).expect("lossless counter encoding");
929            assert_eq!(metrics.replicated_count, 0);
930            assert_eq!(metrics.stats["arn:target"].replicated_count, 1000);
931            assert_eq!(metrics.stats["arn:target"].failed.count, 0);
932            assert_eq!(metrics.q_stat.max.count, expected, "encoded: {encoded}");
933        }
934
935        for invalid in [
936            "-1",
937            "-0.0",
938            "1.5",
939            "999999999999999.01",
940            "1e-400",
941            "18446744073709551616.0",
942            "1e400",
943            "NaN",
944            "Infinity",
945            "\"1\"",
946            "null",
947        ] {
948            let payload = format!(
949                r#"{{"Stats":null,"queued":{{
950                    "curr":{{"count":{invalid},"bytes":0}},
951                    "avg":{{"count":0,"bytes":0}},
952                    "peak":{{"count":0,"bytes":0}}}}}}"#
953            );
954            assert!(
955                serde_json::from_str::<ReplicationMetrics>(&payload).is_err(),
956                "accepted invalid counter {invalid}"
957            );
958        }
959    }
960
961    #[test]
962    fn metrics_reject_missing_core_mixed_and_inconsistent_minio_fields() {
963        for payload in [
964            r#"{"queued":{"curr":{"count":0,"bytes":0},"avg":{"count":0,"bytes":0},"peak":{"count":0,"bytes":0}}}"#,
965            r#"{"Stats":null}"#,
966            r#"{"Stats":null,"stats":{},"queued":{"curr":{"count":0,"bytes":0},"avg":{"count":0,"bytes":0},"peak":{"count":0,"bytes":0}}}"#,
967            r#"{"Stats":null,"replica_size":0,"queued":{"curr":{"count":0,"bytes":0},"avg":{"count":0,"bytes":0},"peak":{"count":0,"bytes":0}}}"#,
968            r#"{"Stats":null,"queued":{"avg":{"count":0,"bytes":0},"peak":{"count":0,"bytes":0}}}"#,
969            r#"{"Stats":null,"queued":{"curr":{"count":0,"bytes":0},"avg":{"count":0,"bytes":0},"peak":{"count":0,"bytes":0}},"failed":{"lastMinute":{"count":0,"bytes":0},"lastHour":{"count":0,"bytes":0}}}"#,
970            r#"{"Stats":{"arn:target":{"failed":{"lastMinute":{"count":0,"bytes":0},"lastHour":{"count":0,"bytes":0},"totals":{"count":2,"bytes":20}},"failedReplicationCount":1,"failedReplicationSize":20}},"queued":{"curr":{"count":0,"bytes":0},"avg":{"count":0,"bytes":0},"peak":{"count":0,"bytes":0}}}"#,
971        ] {
972            assert!(
973                serde_json::from_str::<ReplicationMetrics>(payload).is_err(),
974                "accepted malformed metrics: {payload}"
975            );
976        }
977    }
978}