From 4609e78fcdee9ae1e0e640c46239d36bff7b26d1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Sun, 16 Aug 2026 01:24:36 +0800 Subject: [PATCH] fix(admin): gate, bound, and null-map the mrf stream MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Second review round: - Authorization: the default stream enumerates object names and version ids, which a metrics-only principal must not see — it now requires admin:ReplicationDiff (MinIO parity, route policy updated); ?aggregate=true carries no object identities and keeps admin:GetReplicationMetrics. - The nil UUID is RustFS's in-memory null-version sentinel and now leaves as the S3 wire token 'null' instead of a zero UUID (a pre-versioning object scanned after versioning + existing-object replication can persist it into the ledger). - The durable ledger is not bounded by the in-memory pending cap and the body is buffered before send; the stream now stops at 10,000 documents and signals truncation via x-rustfs-replication-mrf-truncated (mirroring the diff endpoint) plus a warn event, instead of staging an unbounded body. --- rustfs/src/admin/handlers/replication.rs | 140 ++++++++++++++++++++--- rustfs/src/admin/route_policy.rs | 5 +- 2 files changed, 126 insertions(+), 19 deletions(-) diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 6255d6815..b0c34c5e2 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -1263,33 +1263,58 @@ struct MrfEntryDocument { target_arns: Vec, } +/// Upper bound on the number of documents one stream response emits. The +/// durable ledger is not bounded by the in-memory pending cap (recovery can +/// persist far larger generations), and the body is buffered before send, so +/// an unbounded read could stage hundreds of MB per request. A truncated +/// stream is signalled via `x-rustfs-replication-mrf-truncated`. +const REPLICATION_MRF_MAX_STREAM_ENTRIES: usize = 10_000; + /// Project the durable backlog into madmin `ReplicationMRF` documents, /// scoped to `bucket` when it is non-empty (madmin allows an empty bucket to -/// mean "across all buckets"). +/// mean "across all buckets"), bounded by +/// [`REPLICATION_MRF_MAX_STREAM_ENTRIES`]. Returns the documents and whether +/// the backlog was truncated. fn mrf_entry_documents( bucket: &str, durable: &crate::admin::storage_api::replication::DurableMrfBacklog, -) -> Vec { - durable +) -> (Vec, bool) { + let mut documents = Vec::new(); + let mut truncated = false; + for entry in durable .entries .iter() .filter(|entry| bucket.is_empty() || entry.bucket == bucket) - .map(|entry| MrfEntryDocument { + { + if documents.len() >= REPLICATION_MRF_MAX_STREAM_ENTRIES { + truncated = true; + break; + } + documents.push(MrfEntryDocument { node_name: String::new(), bucket: entry.bucket.clone(), object: entry.object.clone(), // Delete-marker purge entries track the marker version separately; // fall back to it so those rows still carry a version identity. + // The nil UUID is RustFS's in-memory null-version sentinel and + // must leave as the S3 wire token, not a zero UUID. version_id: entry .version_id .or(entry.delete_marker_version_id) - .map(|v| v.to_string()) + .map(|v| { + if v.is_nil() { + rustfs_filemeta::NULL_VERSION_ID.to_string() + } else { + v.to_string() + } + }) .unwrap_or_default(), retry_count: entry.retry_count, size: entry.size, target_arns: entry.target_arns.clone(), - }) - .collect() + }); + } + (documents, truncated) } /// Render the MRF backlog as a response body. @@ -1311,17 +1336,18 @@ fn render_mrf_backlog( response: &MrfResponse, durable: &crate::admin::storage_api::replication::DurableMrfBacklog, aggregate: bool, -) -> Result, serde_json::Error> { +) -> Result<(Vec, bool), serde_json::Error> { if aggregate { - return serde_json::to_vec(response); + return Ok((serde_json::to_vec(response)?, false)); } + let (documents, truncated) = mrf_entry_documents(&response.bucket, durable); let mut data = Vec::new(); - for entry in mrf_entry_documents(&response.bucket, durable) { + for entry in documents { serde_json::to_writer(&mut data, &entry)?; data.push(b'\n'); } - Ok(data) + Ok((data, truncated)) } /// `GET /v3/replication/mrf` @@ -1339,14 +1365,28 @@ fn render_mrf_backlog( /// The madmin `node` parameter is accepted but has no filtering effect: the /// durable ledger is cluster-shared with no per-node attribution, so every /// node serves the same (complete) backlog. +/// +/// Authorization: the stream requires `admin:ReplicationDiff` (it enumerates +/// object names and version ids, MinIO parity); `?aggregate=true` carries no +/// object identities and requires only `admin:GetReplicationMetrics`. pub struct ReplicationMrfHandler {} #[async_trait::async_trait] impl Operation for ReplicationMrfHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { - validate_replication_admin_request(&req, AdminAction::GetReplicationMetricsAction).await?; - let queries = extract_query_params(&req.uri); + let aggregate = queries.get("aggregate").map(String::as_str) == Some("true"); + // The default stream enumerates object names and version ids, which + // a metrics-only principal must not see; gate it on the same action + // MinIO uses for this endpoint. The aggregate counters carry no + // object identities and keep the metrics action. + let action = if aggregate { + AdminAction::GetReplicationMetricsAction + } else { + AdminAction::ReplicationDiff + }; + validate_replication_admin_request(&req, action).await?; + let Some(bucket) = queries.get("bucket").filter(|b| !b.is_empty()).cloned() else { return Err(s3_error!(InvalidRequest, "bucket is required")); }; @@ -1379,7 +1419,6 @@ impl Operation for ReplicationMrfHandler { let durable = crate::admin::storage_api::replication::read_durable_mrf_backlog(store).await; let bucket_stats = cluster_replication_stats(&bucket, app_context_from_req(&req)).await; - let aggregate = queries.get("aggregate").map(String::as_str) == Some("true"); let response = build_mrf_response(bucket, &bucket_stats, &durable); if !durable.available && !aggregate { @@ -1397,10 +1436,21 @@ impl Operation for ReplicationMrfHandler { )); } - let data = render_mrf_backlog(&response, &durable, aggregate) + let (data, truncated) = render_mrf_backlog(&response, &durable, aggregate) .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize failed: {e}")))?; let mut headers = HeaderMap::new(); headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json")); + if truncated { + // The madmin stream has no envelope to carry truncation; signal + // it out-of-band (madmin/mc ignore unknown headers), mirroring + // x-rustfs-replication-diff-truncated. + tracing::warn!( + bucket = %response.bucket, + max_entries = REPLICATION_MRF_MAX_STREAM_ENTRIES, + "replication mrf stream truncated; narrow with ?bucket= or drain the backlog" + ); + headers.insert("x-rustfs-replication-mrf-truncated", HeaderValue::from_static("true")); + } Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers)) } } @@ -1770,7 +1820,7 @@ mod tests { let durable = sample_durable_backlog(); let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable); - let body = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize"); + let (body, _) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize"); let text = String::from_utf8(body).expect("body should be utf-8"); let lines: Vec<&str> = text.lines().filter(|line| !line.trim().is_empty()).collect(); @@ -1801,7 +1851,7 @@ mod tests { }; let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable); - let body = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize"); + let (body, _) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize"); assert!( body.is_empty(), "empty backlog must serialize to an empty body, got: {}", @@ -1815,7 +1865,7 @@ mod tests { let durable = sample_durable_backlog(); let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable); - let body = render_mrf_backlog(&response, &durable, true).expect("aggregate body should serialize"); + let (body, _) = render_mrf_backlog(&response, &durable, true).expect("aggregate body should serialize"); let json: serde_json::Value = serde_json::from_slice(&body).expect("aggregate body should be one JSON object"); assert_eq!(json["Bucket"], "bucket-a"); assert_eq!(json["DurableCount"], 1); @@ -1825,6 +1875,60 @@ mod tests { assert_eq!(json["PerObjectEntriesAvailable"], true); } + /// The nil UUID is RustFS's in-memory null-version sentinel; the wire + /// token is `null`, never the zero UUID (second review round). + #[test] + fn mrf_stream_maps_nil_version_to_null_token() { + let durable = DurableMrfBacklog { + available: true, + entries: vec![MrfReplicateEntry { + bucket: "bucket-a".to_string(), + object: "null-version-object".to_string(), + version_id: Some(uuid::Uuid::nil()), + retry_count: 1, + size: 10, + op: MrfOpKind::Object, + ..Default::default() + }], + }; + let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable); + + let (body, truncated) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize"); + assert!(!truncated); + let doc: serde_json::Value = + serde_json::from_str(String::from_utf8(body).expect("utf-8").lines().next().expect("one line")) + .expect("line should be a JSON document"); + assert_eq!(doc["versionId"], "null"); + } + + /// The durable ledger is not bounded by the in-memory pending cap; the + /// stream must stop at the documented bound and signal truncation + /// (second review round). + #[test] + fn mrf_stream_truncates_at_the_documented_bound() { + let entries = (0..super::REPLICATION_MRF_MAX_STREAM_ENTRIES + 1) + .map(|index| MrfReplicateEntry { + bucket: "bucket-a".to_string(), + object: format!("object-{index}"), + retry_count: 1, + op: MrfOpKind::Object, + ..Default::default() + }) + .collect(); + let durable = DurableMrfBacklog { + available: true, + entries, + }; + let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable); + + let (body, truncated) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize"); + assert!(truncated, "one entry past the bound must signal truncation"); + assert_eq!( + String::from_utf8(body).expect("utf-8").lines().count(), + super::REPLICATION_MRF_MAX_STREAM_ENTRIES + ); + } + #[test] fn test_extract_query_params_decodes_percent_encoded_values() { let uri: Uri = "/rustfs/admin/v3/list-remote-targets?bucket=foo%2Fbar&flag=a+b" diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index 774ec7175..474d3ad0f 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -1459,10 +1459,13 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ REPLICATION_DIFF, RouteRiskLevel::Sensitive, ), + // The default stream enumerates object names/version ids and requires + // ReplicationDiff (MinIO parity); only ?aggregate=true relaxes to + // GetReplicationMetrics in the handler. admin( HttpMethod::Get, "/rustfs/admin/v3/replication/mrf", - GET_REPLICATION_METRICS, + REPLICATION_DIFF, RouteRiskLevel::Sensitive, ), ];