Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
577 changes: 565 additions & 12 deletions crates/cli/src/commands/replicate.rs

Large diffs are not rendered by default.

28 changes: 28 additions & 0 deletions crates/cli/tests/fixtures/output_v3/replication/mrf.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
{
"schema_version": 3,
"type": "replication",
"status": "success",
"data": {
"operation": "mrf",
"bucket": "source",
"runtime_stats_available": true,
"cluster": {
"state": "partial",
"observed_nodes": 1,
"expected_nodes": 2
},
"failed": {"count": 2, "size_bytes": 20},
"queue": {"count": 1, "size_bytes": 10},
"durable_backlog": {"available": true, "count": 3, "size_bytes": 30},
"per_object_entries_available": false,
"per_target_durable_entries_available": false,
"targets": [{
"target_arn": "arn:rustfs:replication::target",
"failed_count": 2,
"failed_size_bytes": 20,
"observation_scope": "partial_cluster",
"extensions": {}
}],
"extensions": {}
}
}
46 changes: 46 additions & 0 deletions crates/cli/tests/fixtures/output_v3/replication/status.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
{
"schema_version": 3,
"type": "replication",
"status": "success",
"data": {
"operation": "status",
"bucket": "source",
"availability": "available",
"cluster": {
"state": "partial",
"observed_nodes": 1,
"expected_nodes": 2
},
"queue": {
"count": 1,
"size_bytes": 10,
"scope": "partial_cluster"
},
"totals": {
"replica_count": 4,
"replica_size_bytes": 40,
"replicated_count": 3,
"replicated_size_bytes": 30
},
"targets": [{
"target_arn": "arn:rustfs:replication::target",
"replicated_count": 3,
"replicated_size_bytes": 30,
"failed_count": 1,
"failed_size_bytes": 10,
"latency": {
"average_ms": 1.0,
"current_ms": 2.0,
"maximum_ms": 3.0,
"scope": "partial_cluster"
},
"bandwidth": {
"limit_bytes_per_sec": 100,
"current_bytes_per_sec": 50.0,
"scope": "node_local"
},
"extensions": {}
}],
"extensions": {}
}
}
10 changes: 10 additions & 0 deletions crates/cli/tests/output_schema_v3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,16 @@ fn every_v3_family_has_valid_contract_fixtures() {
}
}

#[test]
fn replication_status_and_mrf_fixtures_are_valid() {
let validator = load_validator(3);
for case in ["status", "mrf"] {
let path = fixture_path("replication", case);
let value = load_json(&path);
assert_valid(&validator, &value, &path.display().to_string());
}
}

#[test]
fn multipart_partial_fixture_preserves_successes_and_per_upload_errors() {
let validator = load_validator(3);
Expand Down
9 changes: 5 additions & 4 deletions crates/core/src/admin/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,11 @@ pub use oidc::{
OidcValidationResult,
};
pub use replication::{
MAX_REPLICATION_DIFF_RESPONSE_BYTES, ReplicationDiff, ReplicationDiffApi, ReplicationDiffEntry,
MAX_REPLICATION_DIFF_RESPONSE_BYTES, MAX_REPLICATION_INSPECTION_RESPONSE_BYTES,
ReplicationCountSize, ReplicationDiff, ReplicationDiffApi, ReplicationDiffEntry,
ReplicationInspectionApi, ReplicationLatencyMetric, ReplicationMetricScope, ReplicationMetrics,
ReplicationMrf, ReplicationMrfTarget, ReplicationQueueMetric, ReplicationTargetMetric,
ReplicationTransferRate,
};
pub use site::{
MAX_SITE_REPLICATION_CA_CERT_BYTES, MAX_SITE_REPLICATION_ERROR_RESPONSE_BYTES,
Expand Down Expand Up @@ -324,9 +328,6 @@ pub trait AdminApi: Send + Sync {
/// Remove a remote replication target
async fn remove_remote_target(&self, bucket: &str, arn: &str) -> Result<()>;

/// Get replication metrics for a bucket
async fn replication_metrics(&self, bucket: &str) -> Result<serde_json::Value>;

// ==================== Service Control Operations ====================

/// Request a service action (restart, stop, freeze, unfreeze)
Expand Down
207 changes: 206 additions & 1 deletion crates/core/src/admin/replication.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
//! Typed contracts for read-only replication diff inspection.
//! Typed contracts for read-only replication inspection.

use std::collections::BTreeMap;

Expand All @@ -11,6 +11,189 @@ use crate::Result;

/// Maximum encoded size accepted for one replication diff response.
pub const MAX_REPLICATION_DIFF_RESPONSE_BYTES: usize = 8 * 1024 * 1024;
/// Maximum encoded size accepted for metrics and MRF responses.
pub const MAX_REPLICATION_INSPECTION_RESPONSE_BYTES: usize = 8 * 1024 * 1024;

/// Scope of a replication observation. Unknown values are retained for forward compatibility.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReplicationMetricScope {
Unavailable,
NodeLocal,
ClusterAggregated,
PartialCluster,
Unknown(String),
}

impl ReplicationMetricScope {
pub fn as_str(&self) -> &str {
match self {
Self::Unavailable => "unavailable",
Self::NodeLocal => "node_local",
Self::ClusterAggregated => "cluster_aggregated",
Self::PartialCluster => "partial_cluster",
Self::Unknown(value) => value,
}
}
}

impl Serialize for ReplicationMetricScope {
fn serialize<S: serde::Serializer>(
&self,
serializer: S,
) -> std::result::Result<S::Ok, S::Error> {
serializer.serialize_str(self.as_str())
}
}

impl<'de> Deserialize<'de> for ReplicationMetricScope {
fn deserialize<D: serde::Deserializer<'de>>(
deserializer: D,
) -> std::result::Result<Self, D::Error> {
let value = String::deserialize(deserializer)?;
Ok(match value.as_str() {
"unavailable" => Self::Unavailable,
"node_local" => Self::NodeLocal,
"cluster_aggregated" => Self::ClusterAggregated,
"partial_cluster" => Self::PartialCluster,
_ => Self::Unknown(value),
})
}
}

#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ReplicationCountSize {
pub count: u64,
#[serde(rename = "bytes", alias = "size")]
pub size: u64,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ReplicationQueueMetric {
pub curr: ReplicationCountSize,
pub avg: ReplicationCountSize,
pub max: ReplicationCountSize,
pub last_minute: ReplicationCountSize,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ReplicationLatencyMetric {
pub avg: f64,
pub curr: f64,
pub max: f64,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct ReplicationTransferRate {
pub avg: f64,
pub curr: f64,
pub peak: f64,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ReplicationTargetMetric {
pub replicated_size: u64,
pub replicated_count: u64,
pub failed: ReplicationCountSize,
#[serde(default)]
pub fail_stats: Option<ReplicationCountSize>,
pub latency: ReplicationLatencyMetric,
pub xfer_rate_lrg: ReplicationTransferRate,
pub xfer_rate_sml: ReplicationTransferRate,
pub bandwidth_limit_bytes_per_sec: u64,
pub current_bandwidth_bytes_per_sec: f64,
#[serde(default)]
pub latency_scope: Option<ReplicationMetricScope>,
#[serde(default)]
pub bandwidth_scope: Option<ReplicationMetricScope>,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ReplicationMetrics {
pub stats: BTreeMap<String, ReplicationTargetMetric>,
pub replica_size: u64,
pub replica_count: u64,
pub replicated_size: u64,
pub replicated_count: u64,
pub q_stat: ReplicationQueueMetric,
#[serde(default)]
pub provider_available: Option<bool>,
#[serde(default)]
pub cluster_complete: Option<bool>,
#[serde(default)]
pub observed_node_count: Option<u32>,
#[serde(default)]
pub expected_node_count: Option<u32>,
#[serde(default)]
pub queue_scope: Option<ReplicationMetricScope>,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ReplicationMrfTarget {
#[serde(rename = "ARN")]
pub arn: String,
#[serde(rename = "FailedCount")]
pub failed_count: u64,
#[serde(rename = "FailedSize")]
pub failed_size: u64,
#[serde(rename = "ObservationScope")]
pub observation_scope: ReplicationMetricScope,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ReplicationMrf {
#[serde(rename = "Bucket")]
pub bucket: String,
#[serde(rename = "Targets")]
pub targets: Vec<ReplicationMrfTarget>,
#[serde(rename = "TotalFailedCount")]
pub total_failed_count: u64,
#[serde(rename = "TotalFailedSize")]
pub total_failed_size: u64,
#[serde(rename = "QueuedCount")]
pub queued_count: u64,
#[serde(rename = "QueuedSize")]
pub queued_size: u64,
#[serde(rename = "PerObjectEntriesAvailable")]
pub per_object_entries_available: bool,
#[serde(rename = "RuntimeStatsAvailable")]
pub runtime_stats_available: bool,
#[serde(rename = "ClusterComplete")]
pub cluster_complete: bool,
#[serde(rename = "ObservedNodeCount")]
pub observed_node_count: u32,
#[serde(rename = "ExpectedNodeCount")]
pub expected_node_count: u32,
#[serde(rename = "DurableBacklogAvailable")]
pub durable_backlog_available: bool,
#[serde(rename = "DurableCount")]
pub durable_count: u64,
#[serde(rename = "DurableSize")]
pub durable_size: u64,
#[serde(rename = "PerTargetDurableEntriesAvailable")]
pub per_target_durable_entries_available: bool,
#[serde(flatten, default)]
pub extra: BTreeMap<String, Value>,
}

#[async_trait]
pub trait ReplicationInspectionApi: Send + Sync {
async fn replication_metrics(&self, bucket: &str) -> Result<ReplicationMetrics>;
async fn replication_mrf(&self, bucket: &str) -> Result<ReplicationMrf>;
}

/// A bounded, on-demand scan of object versions that have not replicated.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
Expand Down Expand Up @@ -120,4 +303,26 @@ mod tests {
let payload = r#"{"Entries":[]}"#;
assert!(serde_json::from_str::<ReplicationDiff>(payload).is_err());
}

#[test]
fn metrics_distinguish_legacy_metadata_and_preserve_unknown_scope() {
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");
assert_eq!(legacy.provider_available, None);
assert_eq!(legacy.cluster_complete, None);

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");
assert_eq!(current.provider_available, Some(true));
assert_eq!(
current.queue_scope,
Some(ReplicationMetricScope::Unknown("future_scope".into()))
);
}

#[test]
fn metrics_and_mrf_reject_negative_counters() {
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}}}"#;
assert!(serde_json::from_str::<ReplicationMetrics>(metrics).is_err());
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}"#;
assert!(serde_json::from_str::<ReplicationMrf>(mrf).is_err());
}
}
Loading