Compare commits

...

4 Commits

Author SHA1 Message Date
唐小鸭 111d10027a fix(admin): serialize replication metrics in minio-go wire shapes
?replication-metrics[=2] and the admin replicationmetrics endpoint
serialized the internal snake_case BucketStats family straight onto the
wire, so 'mc replicate status' decoded all zeros without any error
(backlog#1675 P1-11). The internal structs cannot be renamed: they are
the intra-cluster peer-RPC wire format (rmp_serde to_vec_named in
node_service.rs), pinned by a new regression test.

- New admin/replication_metrics_wire.rs: Serialize-only projections onto
  minio-go replication.Metrics (v1 body, currStats) and MetricsV2
  (uptime/currStats/queueStats/downtimeInfo) with the exact json tags;
  per-target failed becomes the TimedErrStats envelope fed from the
  FailStats rolling window; the queue peak is dual-emitted as max
  (MinIO server tag) and peak (minio-go tag).
- queueStats synthesizes one node from the bucket queue snapshot — the
  aggregation path leaves queue_stats.nodes empty, and mc treats an
  empty node list as 'no data' — and carries transfer summaries
  (Large/Small/Total) derived from the per-target xfer rates.
- Both endpoints share the DTOs; source-health extension keys
  (provider_available/cluster_complete/...) ride along and are ignored
  by Go decoders.
- Widen the ecstore replication_stats_boundary re-exports
  (BucketReplicationStat/InQueueMetric/XferStats) so the admin facade
  chain can name the projected types.
2026-08-15 09:12:56 +08:00
唐小鸭 34b4cfb466 test(admin): pin minio-go Metrics/MetricsV2 wire contract for replication metrics
Red-light evidence for backlog#1675 P1-11: ?replication-metrics[=2]
serializes the internal snake_case BucketStats family straight onto the
wire, while minio-go's replication.Metrics/MetricsV2 expect camelCase
tags (currStats/queueStats/replicaCount/queued/...). Go's decoder is
case-insensitive but does not ignore underscores, so 'mc replicate
status' shows all zeros without any error. The rewritten snapshot tests
assert the minio-go tags (plus a synthesized queueStats node — the
aggregation path leaves queue_stats.nodes empty today) and fail against
the current pass-through serialization.
2026-08-15 08:41:15 +08:00
Zhengchao An 72fd7339c9 test(utils): allow ephemeral port reuse (#6122)
* test(utils): allow ephemeral port reuse

* test(kms): allow any ciphertext prefix
2026-08-15 08:32:10 +08:00
Zhengchao An 71e83aeec4 fix(ci): pin Docker images to release source (#6121) 2026-08-15 07:13:37 +08:00
12 changed files with 640 additions and 36 deletions
+25 -1
View File
@@ -94,6 +94,7 @@ jobs:
short_sha: ${{ steps.check.outputs.short_sha }}
is_prerelease: ${{ steps.check.outputs.is_prerelease }}
create_latest: ${{ steps.check.outputs.create_latest }}
source_ref: ${{ steps.check.outputs.source_ref }}
steps:
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
@@ -118,6 +119,7 @@ jobs:
short_sha=""
is_prerelease=false
create_latest=false
source_ref="$GITHUB_SHA"
if [[ "${{ github.event_name }}" == "workflow_run" ]]; then
# Triggered by build workflow completion
@@ -137,6 +139,7 @@ jobs:
# Extract version info from commit message or use commit SHA
# Use Git to generate consistent short SHA (ensures uniqueness like build.yml)
short_sha=$(git rev-parse --short "$HEAD_SHA")
source_ref="$HEAD_SHA"
# Determine build type based on triggering workflow event and ref
triggering_event="$TRIGGERING_EVENT"
@@ -261,6 +264,23 @@ jobs:
echo "⚠️ Only release versions (latest, v1.0.0, 1.0.0) and prereleases (v1.0.0-alpha1, 1.0.0-beta2) are supported"
;;
esac
if [[ "$should_build" == true && "$input_version" != "latest" ]]; then
tag_ref="refs/tags/$input_version"
if ! git ls-remote --exit-code origin "$tag_ref" >/dev/null 2>&1; then
if [[ "$input_version" == v* ]]; then
tag_ref="refs/tags/${input_version#v}"
else
tag_ref="refs/tags/v$input_version"
fi
fi
if ! git ls-remote --exit-code origin "$tag_ref" >/dev/null 2>&1; then
echo "❌ Release tag not found for Docker build: $input_version"
exit 1
fi
source_ref="$tag_ref"
fi
fi
{
@@ -271,6 +291,7 @@ jobs:
echo "short_sha=$short_sha"
echo "is_prerelease=$is_prerelease"
echo "create_latest=$create_latest"
echo "source_ref=$source_ref"
} >> "$GITHUB_OUTPUT"
echo "🐳 Docker Build Summary:"
@@ -281,6 +302,7 @@ jobs:
echo " - Short SHA: $short_sha"
echo " - Is prerelease: $is_prerelease"
echo " - Create latest: $create_latest"
echo " - Source ref: $source_ref"
# Build multi-arch Docker images
# Strategy: Build images using pre-built binaries from dl.rustfs.com
@@ -308,6 +330,7 @@ jobs:
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: ${{ needs.build-check.outputs.source_ref }}
- name: Login to Docker Hub
uses: docker/login-action@c94ce9fb468520275223c153574b00df6fe4bcc9 # v3
@@ -397,7 +420,8 @@ jobs:
LABELS="org.opencontainers.image.title=RustFS"
LABELS="$LABELS,org.opencontainers.image.description=RustFS distributed object storage system"
LABELS="$LABELS,org.opencontainers.image.version=$VERSION"
LABELS="$LABELS,org.opencontainers.image.revision=${{ github.sha }}"
SOURCE_REVISION="$(git rev-parse HEAD)"
LABELS="$LABELS,org.opencontainers.image.revision=$SOURCE_REVISION"
LABELS="$LABELS,org.opencontainers.image.source=${{ github.server_url }}/${{ github.repository }}"
LABELS="$LABELS,org.opencontainers.image.created=$(date -u +'%Y-%m-%dT%H:%M:%SZ')"
LABELS="$LABELS,org.opencontainers.image.build-type=$BUILD_TYPE"
+12 -11
View File
@@ -184,17 +184,18 @@ pub mod bucket {
mrf_backlog_observability_snapshot,
};
pub use crate::bucket::replication::{
BucketReplicationResyncStatus, BucketReplicationStats, BucketStats, DeleteReplicationConfigSnapshot,
DeletedObjectReplicationInfo, DurableMrfBacklog, DynReplicationPool, MrfOpKind, MrfReplicateEntry,
MustReplicateOptions, ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS,
REMOTE_TARGET_WRITABLE_FIELDS, REPLICATE_INCOMING_DELETE, REPLICATION_CAPABILITY_CONTRACT_VERSION,
REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS, ReplicateDecision, ReplicateObjectInfo,
ReplicationBatchAdmission, ReplicationConfig, ReplicationConfigStructureError, ReplicationConfigurationExt,
ReplicationDeleteScheduleInput, ReplicationDeleteStateSource, ReplicationHealQueueResult, ReplicationObjectBridge,
ReplicationObjectIO, ReplicationOperation, ReplicationPoolTrait, ReplicationPriority, ReplicationQueueAdmission,
ReplicationScannerBridge, ReplicationState, ReplicationStats, ReplicationStatusType, ReplicationStorage,
ReplicationTargetValidationError, ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog,
TargetReplicationResyncStatus, VersionPurgeStatusType, commit_force_delete_intent, complete_force_delete_intent,
BucketReplicationResyncStatus, BucketReplicationStat, BucketReplicationStats, BucketStats,
DeleteReplicationConfigSnapshot, DeletedObjectReplicationInfo, DurableMrfBacklog, DynReplicationPool, InQueueMetric,
MrfOpKind, MrfReplicateEntry, MustReplicateOptions, ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION,
REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, REPLICATE_INCOMING_DELETE,
REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS,
ReplicateDecision, ReplicateObjectInfo, ReplicationBatchAdmission, ReplicationConfig,
ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationDeleteScheduleInput,
ReplicationDeleteStateSource, ReplicationHealQueueResult, ReplicationObjectBridge, ReplicationObjectIO,
ReplicationOperation, ReplicationPoolTrait, ReplicationPriority, ReplicationQueueAdmission, ReplicationScannerBridge,
ReplicationState, ReplicationStats, ReplicationStatusType, ReplicationStorage, ReplicationTargetValidationError,
ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog, TargetReplicationResyncStatus,
VersionPurgeStatusType, XferStats, commit_force_delete_intent, complete_force_delete_intent,
delete_replication_state_from_config, delete_replication_version_id, get_global_replication_pool,
get_global_replication_stats, init_background_replication, invalid_replication_config_status_field,
persist_force_delete_intent, read_durable_mrf_backlog, replication_state_to_filemeta, replication_status_to_filemeta,
+1 -1
View File
@@ -81,6 +81,6 @@ pub use replication_queue_boundary::{
pub use replication_resync_boundary::{BucketReplicationResyncStatus, ResyncOpts, TargetReplicationResyncStatus};
pub use replication_scanner_bridge::ReplicationScannerBridge;
pub use replication_state::{ReplicationStats, RuntimeReplicationTargetBacklog};
pub use replication_stats_boundary::{BucketReplicationStats, BucketStats};
pub use replication_stats_boundary::{BucketReplicationStat, BucketReplicationStats, BucketStats, InQueueMetric, XferStats};
pub use replication_storage_boundary::{ReplicationObjectIO, ReplicationStorage};
pub(crate) use replication_target_config_bridge::ReplicationTargetConfigBridge;
@@ -15,7 +15,9 @@
#[cfg(test)]
pub(crate) use rustfs_replication::FailStats;
pub(crate) use rustfs_replication::{
ActiveWorkerStat, BucketReplicationStat, InQueueMetric, ProxyMetric, ProxyStatsCache, QueueCache, ReplicationMetricScope,
SRMetricsSummary, XferStats,
ActiveWorkerStat, ProxyMetric, ProxyStatsCache, QueueCache, ReplicationMetricScope, SRMetricsSummary,
};
pub use rustfs_replication::{BucketReplicationStats, BucketStats};
// Public so the admin wire DTOs (rustfs/src/admin/replication_metrics_wire.rs)
// can project the internal stats onto the minio-go response shapes through
// the storage_api facade chain.
pub use rustfs_replication::{BucketReplicationStat, BucketReplicationStats, BucketStats, InQueueMetric, XferStats};
-2
View File
@@ -225,8 +225,6 @@ async fn nothing_readable_leaves_the_bundle_unwrapped() {
"artifact {} carries the raw on-disk record",
artifact.path
);
// A cheap structural check too: an encrypted payload is not JSON.
assert_ne!(payload.first(), Some(&b'{'), "artifact {} looks like plaintext JSON", artifact.path);
}
// The manifest itself is not encrypted, so assert directly that it carries
-3
View File
@@ -659,9 +659,6 @@ mod test {
// Port should be in valid range (u16 max is always <= 65535)
assert!(port1 > 0);
assert!(port2 > 0);
// Different calls should typically return different ports
assert_ne!(port1, port2);
}
#[test]
+7 -2
View File
@@ -535,8 +535,13 @@ impl Operation for GetReplicationMetricsHandler {
let bucket_stats = cluster_replication_stats(bucket, app_context_from_req(&req)).await;
let data = serde_json::to_vec(&bucket_stats.replication_stats)
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "serialize failed"))?;
// Same minio-go `replication.Metrics` wire shape as
// `?replication-metrics` — the internal snake_case stats are the peer
// RPC wire format and must not leak here.
let data = serde_json::to_vec(&crate::admin::replication_metrics_wire::MetricsWire::from(
&bucket_stats.replication_stats,
))
.map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "serialize failed"))?;
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
+1
View File
@@ -19,6 +19,7 @@ pub mod handlers;
mod plugin_contract;
// Contract inventory is validated by tests before later runtime integration.
#[allow(dead_code)]
pub(crate) mod replication_metrics_wire;
pub(crate) mod route_policy;
pub mod router;
pub(crate) mod runtime_sources;
@@ -0,0 +1,511 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Serialize-only wire projections of the internal replication statistics
//! onto the minio-go `replication.Metrics` / `replication.MetricsV2` json
//! shapes consumed by `mc replicate status` (`?replication-metrics[=2]` and
//! the admin `replicationmetrics` endpoint).
//!
//! Red line: the internal `BucketStats` family in
//! `crates/replication/src/stats.rs` is ALSO the intra-cluster peer-RPC wire
//! format — `node_service.rs` encodes it with `rmp_serde::to_vec_named`, so
//! its Rust field names travel between nodes as msgpack map keys. Renaming
//! those serde names would break mixed-version clusters mid rolling upgrade.
//! All madmin/minio-go interop therefore happens in these DTOs; never add
//! `#[serde(rename)]` to the internal structs instead.
//!
//! Field names below are the exact json tags of minio-go
//! `pkg/replication/replication.go` (v7.0.91). Keys minio-go does not know
//! are RustFS extensions; Go decoders ignore unknown keys. `max`/`peak` are
//! both emitted for the queue peak because the MinIO server writes `max`
//! while minio-go reads `peak` (an upstream drift); emitting both keeps every
//! decoder working.
use serde::Serialize;
use std::collections::HashMap;
use std::time::Duration;
use crate::admin::storage_api::replication::{
BucketReplicationStat as InternalReplicationStat, BucketReplicationStats as InternalReplicationStats, BucketStats,
InQueueMetric as InternalInQueueMetric, XferStats as InternalXferStats,
};
/// minio-go `replication.RStat`.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct RStatWire {
#[serde(rename = "count")]
pub count: f64,
#[serde(rename = "bytes")]
pub bytes: i64,
}
/// minio-go `replication.TimedErrStats`.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct TimedErrStatsWire {
#[serde(rename = "lastMinute")]
pub last_minute: RStatWire,
#[serde(rename = "lastHour")]
pub last_hour: RStatWire,
#[serde(rename = "totals")]
pub totals: RStatWire,
}
impl TimedErrStatsWire {
fn add(self, other: TimedErrStatsWire) -> TimedErrStatsWire {
fn add(a: RStatWire, b: RStatWire) -> RStatWire {
RStatWire {
count: a.count + b.count,
bytes: a.bytes.saturating_add(b.bytes),
}
}
TimedErrStatsWire {
last_minute: add(self.last_minute, other.last_minute),
last_hour: add(self.last_hour, other.last_hour),
totals: add(self.totals, other.totals),
}
}
}
/// minio-go `replication.QStat`.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct QStatWire {
#[serde(rename = "count")]
pub count: f64,
#[serde(rename = "bytes")]
pub bytes: f64,
}
/// minio-go `replication.InQueueMetric`, with the queue peak emitted under
/// both `peak` (minio-go tag) and `max` (MinIO server tag).
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct InQueueMetricWire {
#[serde(rename = "curr")]
pub curr: QStatWire,
#[serde(rename = "avg")]
pub avg: QStatWire,
#[serde(rename = "max")]
pub max: QStatWire,
#[serde(rename = "peak")]
pub peak: QStatWire,
}
impl From<&InternalInQueueMetric> for InQueueMetricWire {
fn from(metric: &InternalInQueueMetric) -> Self {
fn qstat(bytes: i64, count: i64) -> QStatWire {
QStatWire {
count: count as f64,
bytes: bytes as f64,
}
}
let peak = qstat(metric.max.bytes, metric.max.count);
InQueueMetricWire {
curr: qstat(metric.curr.bytes, metric.curr.count),
avg: qstat(metric.avg.bytes, metric.avg.count),
max: peak,
peak,
}
}
}
/// minio-go `replication.XferStats`.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct XferStatsWire {
#[serde(rename = "avgRate")]
pub avg_rate: f64,
#[serde(rename = "peakRate")]
pub peak_rate: f64,
#[serde(rename = "currRate")]
pub curr_rate: f64,
}
impl XferStatsWire {
fn merge(self, other: XferStatsWire) -> XferStatsWire {
XferStatsWire {
avg_rate: self.avg_rate + other.avg_rate,
peak_rate: self.peak_rate.max(other.peak_rate),
curr_rate: self.curr_rate + other.curr_rate,
}
}
}
impl From<&InternalXferStats> for XferStatsWire {
fn from(stats: &InternalXferStats) -> Self {
XferStatsWire {
avg_rate: stats.avg,
peak_rate: stats.peak,
curr_rate: stats.curr,
}
}
}
/// minio-go `replication.WorkerStat`. RustFS does not track per-bucket worker
/// occupancy yet, so this always reports zeros.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct WorkerStatWire {
#[serde(rename = "curr")]
pub curr: i32,
#[serde(rename = "avg")]
pub avg: f32,
#[serde(rename = "max")]
pub max: i32,
}
/// minio-go `replication.ReplMRFStats`. RustFS does not track the 5-minute /
/// dropped MRF windows, so this always reports zeros; the durable backlog is
/// enumerable via `/v3/replication/mrf` instead.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct ReplMrfStatsWire {
#[serde(rename = "failedCount_last5min")]
pub last_failed_count: u64,
#[serde(rename = "droppedCount_since_uptime")]
pub total_dropped_count: u64,
#[serde(rename = "droppedBytes_since_uptime")]
pub total_dropped_bytes: u64,
}
/// minio-go `replication.CounterSummary`.
#[derive(Debug, Default, Clone, Copy, Serialize)]
pub(crate) struct CounterSummaryWire {
#[serde(rename = "last1hr")]
pub last1hr: u64,
#[serde(rename = "last1m")]
pub last1m: u64,
#[serde(rename = "total")]
pub total: u64,
}
/// minio-go `replication.TargetMetrics` (one remote target / ARN).
#[derive(Debug, Default, Serialize)]
pub(crate) struct TargetMetricsWire {
#[serde(rename = "replicationCount")]
pub replicated_count: i64,
#[serde(rename = "completedReplicationSize")]
pub replicated_size: i64,
/// Bandwidth limit for this target. The tag says "bits" but both MinIO
/// and minio-go treat the value as bytes/sec; keep bytes/sec.
#[serde(rename = "limitInBits")]
pub bandwidth_limit_bytes_per_sec: i64,
#[serde(rename = "currentBandwidth")]
pub current_bandwidth_bytes_per_sec: f64,
#[serde(rename = "failed")]
pub failed: TimedErrStatsWire,
#[serde(rename = "failedReplicationSize")]
pub failed_size: i64,
#[serde(rename = "failedReplicationCount")]
pub failed_count: i64,
}
fn target_timed_err_stats(stat: &InternalReplicationStat) -> TimedErrStatsWire {
let last_minute = stat.fail_stats.recent_since(Duration::from_secs(60));
let last_hour = stat.fail_stats.recent_since(Duration::from_secs(3600));
TimedErrStatsWire {
last_minute: RStatWire {
count: last_minute.count as f64,
bytes: last_minute.size,
},
last_hour: RStatWire {
count: last_hour.count as f64,
bytes: last_hour.size,
},
totals: RStatWire {
count: stat.failed.count as f64,
bytes: stat.failed.size,
},
}
}
impl From<&InternalReplicationStat> for TargetMetricsWire {
fn from(stat: &InternalReplicationStat) -> Self {
TargetMetricsWire {
replicated_count: stat.replicated_count,
replicated_size: stat.replicated_size,
bandwidth_limit_bytes_per_sec: stat.bandwidth_limit_bytes_per_sec,
current_bandwidth_bytes_per_sec: stat.current_bandwidth_bytes_per_sec,
failed: target_timed_err_stats(stat),
failed_size: stat.failed.size,
failed_count: stat.failed.count,
}
}
}
/// minio-go `replication.Metrics` — the `currStats` member of `MetricsV2` and
/// the whole v1 response body. The trailing snake_case fields are RustFS
/// source-health extension keys (ignored by Go decoders) carried over from
/// the previous response shape.
#[derive(Debug, Default, Serialize)]
pub(crate) struct MetricsWire {
#[serde(rename = "Stats")]
pub stats: HashMap<String, TargetMetricsWire>,
#[serde(rename = "completedReplicationSize")]
pub replicated_size: i64,
#[serde(rename = "replicaSize")]
pub replica_size: i64,
#[serde(rename = "replicaCount")]
pub replica_count: i64,
#[serde(rename = "replicationCount")]
pub replicated_count: i64,
#[serde(rename = "failed")]
pub failed: TimedErrStatsWire,
#[serde(rename = "queued")]
pub queued: InQueueMetricWire,
// RustFS extension keys (source health of the aggregation).
pub provider_available: bool,
pub cluster_complete: bool,
pub observed_node_count: u32,
pub expected_node_count: u32,
}
impl From<&InternalReplicationStats> for MetricsWire {
fn from(stats: &InternalReplicationStats) -> Self {
let mut failed = TimedErrStatsWire::default();
let mut targets = HashMap::with_capacity(stats.stats.len());
for (arn, stat) in &stats.stats {
let target = TargetMetricsWire::from(stat);
failed = failed.add(target.failed);
targets.insert(arn.clone(), target);
}
MetricsWire {
stats: targets,
replicated_size: stats.replicated_size,
replica_size: stats.replica_size,
replica_count: stats.replica_count,
replicated_count: stats.replicated_count,
failed,
queued: InQueueMetricWire::from(&stats.q_stat),
provider_available: stats.provider_available,
cluster_complete: stats.cluster_complete,
observed_node_count: stats.observed_node_count,
expected_node_count: stats.expected_node_count,
}
}
}
/// minio-go `replication.ReplQNodeStats`.
#[derive(Debug, Default, Serialize)]
pub(crate) struct ReplQNodeStatsWire {
#[serde(rename = "nodeName")]
pub node_name: String,
#[serde(rename = "uptime")]
pub uptime: i64,
#[serde(rename = "activeWorkers")]
pub workers: WorkerStatWire,
#[serde(rename = "transferSummary")]
pub xfer_stats: XferSummaryWire,
#[serde(rename = "tgtTransferStats")]
pub tgt_xfer_stats: TargetXferSummaryWire,
#[serde(rename = "queueStats")]
pub q_stats: InQueueMetricWire,
#[serde(rename = "mrfStats")]
pub mrf_stats: ReplMrfStatsWire,
#[serde(rename = "retries")]
pub retries: CounterSummaryWire,
#[serde(rename = "errors")]
pub errors: CounterSummaryWire,
}
/// minio-go `replication.ReplQueueStats`.
#[derive(Debug, Default, Serialize)]
pub(crate) struct ReplQueueStatsWire {
#[serde(rename = "nodes")]
pub nodes: Vec<ReplQNodeStatsWire>,
}
/// minio-go `replication.MetricsV2` — the `?replication-metrics=2` body.
#[derive(Debug, Default, Serialize)]
pub(crate) struct MetricsV2Wire {
#[serde(rename = "uptime")]
pub uptime: i64,
#[serde(rename = "currStats")]
pub current_stats: MetricsWire,
#[serde(rename = "queueStats")]
pub queue_stats: ReplQueueStatsWire,
#[serde(rename = "downtimeInfo")]
pub downtime_info: HashMap<String, serde_json::Value>,
}
/// `transferSummary` map keyed by minio-go `MetricName` (Large/Small/Total).
type XferSummaryWire = HashMap<&'static str, XferStatsWire>;
/// `tgtTransferStats` map keyed by target ARN.
type TargetXferSummaryWire = HashMap<String, XferSummaryWire>;
fn transfer_summaries(stats: &InternalReplicationStats) -> (XferSummaryWire, TargetXferSummaryWire) {
let mut summary: XferSummaryWire = HashMap::new();
let mut per_target: TargetXferSummaryWire = HashMap::new();
for (arn, stat) in &stats.stats {
let large = XferStatsWire::from(&stat.xfer_rate_lrg);
let small = XferStatsWire::from(&stat.xfer_rate_sml);
let total = large.merge(small);
per_target.insert(arn.clone(), HashMap::from([("Large", large), ("Small", small), ("Total", total)]));
for (key, value) in [("Large", large), ("Small", small), ("Total", total)] {
let entry = summary.entry(key).or_default();
*entry = entry.merge(value);
}
}
(summary, per_target)
}
impl MetricsV2Wire {
/// Project the aggregated internal stats onto the `MetricsV2` shape.
///
/// The aggregation path leaves `queue_stats.nodes` empty today, so a
/// single node entry is synthesized from the bucket queue snapshot —
/// `mc replicate status` derives its queue/worker panels from
/// `queueStats.nodes` and treats an empty list as "no data".
pub(crate) fn from_stats(bucket_stats: &BucketStats, node_name: &str) -> Self {
let (xfer_stats, tgt_xfer_stats) = transfer_summaries(&bucket_stats.replication_stats);
let mut nodes: Vec<ReplQNodeStatsWire> = bucket_stats
.queue_stats
.nodes
.iter()
.map(|node| ReplQNodeStatsWire {
node_name: node_name.to_string(),
uptime: bucket_stats.uptime,
q_stats: InQueueMetricWire::from(&node.q_stats),
..Default::default()
})
.collect();
if nodes.is_empty() {
nodes.push(ReplQNodeStatsWire {
node_name: node_name.to_string(),
uptime: bucket_stats.uptime,
q_stats: InQueueMetricWire::from(&bucket_stats.replication_stats.q_stat),
xfer_stats: xfer_stats.clone(),
tgt_xfer_stats: tgt_xfer_stats.clone(),
..Default::default()
});
} else {
// Attach the transfer summaries to the first node; the internal
// snapshot does not attribute transfer rates per node.
if let Some(first) = nodes.first_mut() {
first.xfer_stats = xfer_stats.clone();
first.tgt_xfer_stats = tgt_xfer_stats.clone();
}
}
MetricsV2Wire {
uptime: bucket_stats.uptime,
current_stats: MetricsWire::from(&bucket_stats.replication_stats),
queue_stats: ReplQueueStatsWire { nodes },
downtime_info: HashMap::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_bucket_stats() -> BucketStats {
let mut stats = BucketStats {
uptime: 42,
..Default::default()
};
stats.replication_stats.replica_count = 2;
stats.replication_stats.replica_size = 128;
stats.replication_stats.replicated_count = 9;
stats.replication_stats.replicated_size = 4096;
let target = stats
.replication_stats
.stats
.entry("arn:minio:replication::t:b".to_string())
.or_default();
target.replicated_count = 9;
target.replicated_size = 4096;
target.failed.count = 3;
target.failed.size = 900;
target.bandwidth_limit_bytes_per_sec = 1024;
target.current_bandwidth_bytes_per_sec = 512.5;
stats
.replication_stats
.q_stat
.curr
.now_count
.store(4, std::sync::atomic::Ordering::Relaxed);
stats
.replication_stats
.q_stat
.curr
.now_bytes
.store(1200, std::sync::atomic::Ordering::Relaxed);
stats.replication_stats.q_stat = stats.replication_stats.q_stat.snapshot();
stats
}
#[test]
fn metrics_wire_matches_minio_go_tags() {
let stats = sample_bucket_stats();
let json = serde_json::to_value(MetricsWire::from(&stats.replication_stats)).expect("v1 wire should serialize");
assert_eq!(json["replicaCount"], 2);
assert_eq!(json["replicaSize"], 128);
assert_eq!(json["replicationCount"], 9);
assert_eq!(json["completedReplicationSize"], 4096);
assert_eq!(json["queued"]["curr"]["count"], 4.0);
assert_eq!(json["queued"]["curr"]["bytes"], 1200.0);
let target = &json["Stats"]["arn:minio:replication::t:b"];
assert_eq!(target["replicationCount"], 9);
assert_eq!(target["completedReplicationSize"], 4096);
assert_eq!(target["limitInBits"], 1024);
assert_eq!(target["currentBandwidth"], 512.5);
// failed is the madmin TimedErrStats envelope, not the internal
// {count,size} pair.
assert_eq!(target["failed"]["totals"]["count"], 3.0);
assert_eq!(target["failed"]["totals"]["bytes"], 900);
assert!(target["failed"].get("count").is_none());
// Aggregate failed mirrors the per-target totals.
assert_eq!(json["failed"]["totals"]["count"], 3.0);
}
#[test]
fn metrics_v2_wire_synthesizes_queue_node() {
let stats = sample_bucket_stats();
let json = serde_json::to_value(MetricsV2Wire::from_stats(&stats, "node-1:9000")).expect("v2 wire should serialize");
assert_eq!(json["uptime"], 42);
assert_eq!(json["currStats"]["replicaCount"], 2);
let node = &json["queueStats"]["nodes"][0];
assert_eq!(node["nodeName"], "node-1:9000");
assert_eq!(node["uptime"], 42);
assert_eq!(node["queueStats"]["curr"]["count"], 4.0);
// The queue peak is emitted under both the minio-go tag (`peak`) and
// the MinIO server tag (`max`).
assert_eq!(node["queueStats"]["peak"], node["queueStats"]["max"]);
assert!(node["activeWorkers"].get("curr").is_some());
assert!(node["transferSummary"].get("Total").is_some());
assert_eq!(json["downtimeInfo"], serde_json::json!({}));
}
/// Pin the intra-cluster peer-RPC wire format of the internal stats: it
/// is msgpack with the Rust field names as map keys
/// (`rmp_serde::to_vec_named` in node_service.rs). If someone "fixes"
/// the interop bug by renaming the internal serde fields instead of using
/// these DTOs, this test fails and points them here.
#[test]
fn internal_bucket_stats_rpc_wire_stays_snake_case() {
let stats = sample_bucket_stats();
let encoded = rmp_serde::to_vec_named(&stats).expect("internal stats should encode");
let value: serde_json::Value = rmp_serde::from_slice(&encoded).expect("named msgpack should decode generically");
assert!(
value.get("replication_stats").is_some(),
"peer RPC key replication_stats must not be renamed"
);
assert!(value["replication_stats"].get("q_stat").is_some());
assert!(value.get("queue_stats").is_some());
assert!(value.get("proxy_stats").is_some());
let decoded: BucketStats = rmp_serde::from_slice(&encoded).expect("round-trip through the peer RPC wire");
assert_eq!(decoded.replication_stats.replica_count, 2);
}
}
+67 -13
View File
@@ -1548,7 +1548,8 @@ async fn build_replication_metrics_response(
let bucket_stats = apply_replication_metrics_bandwidth_report(bucket_stats, collect_replication_metrics_bandwidth(bucket));
let bucket_stats = apply_replication_metrics_runtime_fields(bucket_stats, route, replication_metrics_uptime_seconds());
let body = serialize_replication_metrics_body(&bucket_stats, route)?;
let node_name = crate::runtime_sources::current_local_node_name().await.unwrap_or_default();
let body = serialize_replication_metrics_body(&bucket_stats, route, &node_name)?;
let mut resp = S3Response::with_status(Body::from(body), StatusCode::OK);
resp.headers
@@ -1608,12 +1609,24 @@ fn apply_replication_metrics_runtime_fields(
bucket_stats
}
fn serialize_replication_metrics_body(bucket_stats: &BucketStats, route: ReplicationExtRoute) -> S3Result<Vec<u8>> {
/// Serialize the metrics body in the minio-go wire shapes
/// (`replication.Metrics` for v1, `replication.MetricsV2` for v2). The
/// internal `BucketStats` serde names are the intra-cluster peer-RPC wire
/// format and must never appear here — see
/// `crate::admin::replication_metrics_wire`.
fn serialize_replication_metrics_body(
bucket_stats: &BucketStats,
route: ReplicationExtRoute,
node_name: &str,
) -> S3Result<Vec<u8>> {
use crate::admin::replication_metrics_wire::{MetricsV2Wire, MetricsWire};
match route {
ReplicationExtRoute::MetricsV1 => {
serde_json::to_vec(&bucket_stats.replication_stats).map_err(|e| s3_error!(InternalError, "{e}"))
serde_json::to_vec(&MetricsWire::from(&bucket_stats.replication_stats)).map_err(|e| s3_error!(InternalError, "{e}"))
}
ReplicationExtRoute::MetricsV2 => {
serde_json::to_vec(&MetricsV2Wire::from_stats(bucket_stats, node_name)).map_err(|e| s3_error!(InternalError, "{e}"))
}
ReplicationExtRoute::MetricsV2 => serde_json::to_vec(bucket_stats).map_err(|e| s3_error!(InternalError, "{e}")),
ReplicationExtRoute::Check | ReplicationExtRoute::ResetStart | ReplicationExtRoute::ResetStatus => {
Err(s3_error!(InternalError, "invalid route for metrics response"))
}
@@ -4147,22 +4160,37 @@ mod tests {
assert!(err.message().unwrap_or_default().contains("rule-stale"));
}
/// The v1 body must decode into minio-go `replication.Metrics` (exact
/// json tags); Go's decoder matches case-insensitively but does not
/// ignore underscores, so the internal snake_case names read as all-zero.
#[test]
fn serialize_replication_metrics_body_v1_returns_replication_stats_only() {
fn serialize_replication_metrics_body_v1_returns_minio_go_metrics_shape() {
let mut stats = BucketStats {
uptime: 99,
..Default::default()
};
stats.replication_stats.replica_count = 7;
stats.replication_stats.replicated_size = 2048;
stats
.replication_stats
.stats
.entry("arn:minio:replication::t:b".to_string())
.or_default()
.replicated_count = 5;
stats.proxy_stats.put_total = 3;
let body =
serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV1).expect("metrics v1 body should serialize");
let body = serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV1, "node-1:9000")
.expect("metrics v1 body should serialize");
let payload: serde_json::Value = serde_json::from_slice(&body).expect("body should be json");
assert_eq!(payload["replica_count"], 7);
assert_eq!(payload["replicaCount"], 7);
assert_eq!(payload["completedReplicationSize"], 2048);
assert_eq!(payload["Stats"]["arn:minio:replication::t:b"]["replicationCount"], 5);
assert!(payload.get("uptime").is_none());
assert!(payload.get("proxy_stats").is_none());
// The internal snake_case names must not leak into the wire body.
assert!(payload.get("replica_count").is_none());
assert!(payload.get("q_stat").is_none());
}
#[test]
@@ -4248,22 +4276,48 @@ mod tests {
assert_eq!(target.current_bandwidth_bytes_per_sec, 3000.0);
}
/// The v2 body must decode into minio-go `replication.MetricsV2`
/// (`uptime`/`currStats`/`queueStats`); `mc replicate status` reads
/// `currStats` and `queueStats.nodes` and silently shows zeros when the
/// keys do not match.
#[test]
fn serialize_replication_metrics_body_v2_returns_full_bucket_stats() {
fn serialize_replication_metrics_body_v2_returns_minio_go_metrics_v2_shape() {
let mut stats = BucketStats {
uptime: 99,
..Default::default()
};
stats.replication_stats.replica_count = 7;
stats
.replication_stats
.q_stat
.curr
.now_count
.store(4, std::sync::atomic::Ordering::Relaxed);
stats
.replication_stats
.q_stat
.curr
.now_bytes
.store(1200, std::sync::atomic::Ordering::Relaxed);
stats.replication_stats.q_stat = stats.replication_stats.q_stat.snapshot();
stats.proxy_stats.put_total = 3;
let body =
serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV2).expect("metrics v2 body should serialize");
let body = serialize_replication_metrics_body(&stats, ReplicationExtRoute::MetricsV2, "node-1:9000")
.expect("metrics v2 body should serialize");
let payload: serde_json::Value = serde_json::from_slice(&body).expect("body should be json");
assert_eq!(payload["uptime"], 99);
assert_eq!(payload["replication_stats"]["replica_count"], 7);
assert_eq!(payload["proxy_stats"]["put_total"], 3);
assert_eq!(payload["currStats"]["replicaCount"], 7);
assert_eq!(payload["currStats"]["queued"]["curr"]["count"], 4.0);
// The queue snapshot must surface at least one node: mc derives the
// worker/queue panels from queueStats.nodes and treats an empty list
// as "no data".
assert_eq!(payload["queueStats"]["nodes"][0]["queueStats"]["curr"]["count"], 4.0);
assert_eq!(payload["queueStats"]["nodes"][0]["uptime"], 99);
// The internal snake_case names must not leak into the wire body.
assert!(payload.get("replication_stats").is_none());
assert!(payload.get("queue_stats").is_none());
assert!(payload.get("proxy_stats").is_none());
}
#[test]
+4
View File
@@ -417,6 +417,10 @@ pub(crate) mod replication {
};
pub(crate) type BucketReplicationResyncStatus = super::ecstore_bucket::replication::BucketReplicationResyncStatus;
pub(crate) type BucketStats = super::ecstore_bucket::replication::BucketStats;
pub(crate) type BucketReplicationStats = super::ecstore_bucket::replication::BucketReplicationStats;
pub(crate) type BucketReplicationStat = super::ecstore_bucket::replication::BucketReplicationStat;
pub(crate) type InQueueMetric = super::ecstore_bucket::replication::InQueueMetric;
pub(crate) type XferStats = super::ecstore_bucket::replication::XferStats;
pub(crate) type ReplicationStatusType = super::ecstore_bucket::replication::ReplicationStatusType;
pub(crate) type ResyncOpts = super::ecstore_bucket::replication::ResyncOpts;
pub(crate) type ResyncStatusType = super::ecstore_bucket::replication::ResyncStatusType;
@@ -195,6 +195,13 @@ IFS= read -r -d '' expected_docker_automatic_guard <<'EOF' || true
EOF
expected_docker_automatic_guard=${expected_docker_automatic_guard%$'\n'}
require_job_if "$docker_workflow" "build-check" "$expected_docker_automatic_guard"
require_line "$docker_workflow" ' source_ref: ${{ steps.check.outputs.source_ref }}' "Docker source ref output"
require_line "$docker_workflow" ' source_ref="$HEAD_SHA"' "automatic Docker source ref"
require_line "$docker_workflow" ' source_ref="$tag_ref"' "manual Docker source ref"
require_line "$docker_workflow" ' ref: ${{ needs.build-check.outputs.source_ref }}' "Docker release source checkout"
require_line "$docker_workflow" ' SOURCE_REVISION="$(git rev-parse HEAD)"' "Docker source revision resolution"
require_line "$docker_workflow" ' LABELS="$LABELS,org.opencontainers.image.revision=$SOURCE_REVISION"' "Docker revision label"
require_absent "$docker_workflow" 'org.opencontainers.image.revision=${{ github.sha }}' "Docker revision must not use the workflow branch SHA"
docker_manual_guard=$(awk '
$0 == " *-preview*)" { in_preview = 1 }