1use 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
12pub const MAX_REPLICATION_DIFF_RESPONSE_BYTES: usize = 8 * 1024 * 1024;
14pub const MAX_REPLICATION_INSPECTION_RESPONSE_BYTES: usize = 8 * 1024 * 1024;
16
17#[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 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#[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#[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#[async_trait]
716pub trait ReplicationDiffApi: Send + Sync {
717 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}