Compare commits

..

3 Commits

Author SHA1 Message Date
唐小鸭 858cf20d9a fix(replication): keep first-replication miss quiet in multipart LWW log
Version-absent is the normal first delivery of a version, not a
degraded comparison skip; only real read errors (quorum loss) warrant
the warn added in the previous commit.
2026-08-21 19:35:22 +08:00
唐小鸭 66e63d20dd fix(replication): surface skipped multipart LWW comparison at warn level
When the complete-multipart LWW gate cannot read the destination
version (quorum error), the inbound metadata is applied without
comparison — the exact overwrite rustfs/backlog#1953 exists to prevent.
Keep the fail-open semantics (failing the complete would loop through
MRF) but log the degraded path at warn so operators can see it; the
version-absent first replication keeps riding the same branch.
2026-08-21 19:33:59 +08:00
唐小鸭 81ef2882df fix(replication): apply receiver-side LWW to inbound metadata categories
Metadata-only replication reuses the whole-object PUT/multipart
transports, so in active-active topologies an inbound authorized
replication write carried the source's tags, retention, and legal hold
verbatim and unconditionally overwrote a category the destination had
modified more recently — both sites ended permanently diverged while
reporting COMPLETED (rustfs/backlog#1953, audit A4/P1-6).

The receiver now judges each category independently under the object
write lock: when the destination version's stored internal timestamp is
strictly newer than the inbound source timestamp, the local category
values and timestamp are kept; the rest of the write proceeds per the
inbound metadata and the object-level result stays successful, so a
local win never feeds an MRF retry loop. When the inbound category wins,
its internal timestamp key is pinned to the source-authored time instead
of the receiver-now() value stamped by the object-lock eval_metadata
path. No stored timestamp (pre-P1-6 data) or no inbound timestamp keeps
today's overwrite behavior.

The put_object hook reuses the existing commit-lock WORM-gate read (no
extra fanout); complete_multipart_upload adds one gated read under its
held lock, and the sender's complete options now carry the three
category timestamps so the multipart transport gets the same receiver
behavior. Local wins restore the timestamp via insert_str so a
MinIO-written single-key version still yields both compatibility keys.
2026-08-21 19:29:26 +08:00
17 changed files with 892 additions and 447 deletions
+5
View File
@@ -61,6 +61,11 @@ mod get_codec_streaming_compat_test;
#[cfg(test)]
mod version_id_regression_test;
// Receiver-side replication LWW (rustfs/backlog#1953): stale inbound
// replication metadata must not overwrite a newer local category state.
#[cfg(test)]
mod replication_lww_receiver_test;
// Data usage regression tests
#[cfg(test)]
mod data_usage_test;
@@ -0,0 +1,152 @@
#![cfg(test)]
// 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.
//! Receiver-side replication LWW over the wire (rustfs/backlog#1953, audit
//! A4/P1-6).
//!
//! In an active-active topology both sites' metadata states arrive at the
//! peer as authorized replication PUTs carrying per-category source
//! timestamps (`x-rustfs-source-replication-tagging-timestamp` header
//! family). Before the fix the receiver applied them unconditionally, so a
//! stale delivery overwrote a newer local state and the two sites diverged
//! permanently while both reported COMPLETED. This test drives one live
//! `rustfs` server with simulated inbound replication PUTs for the same
//! object version and asserts the newer tagging state wins regardless of
//! delivery order, while a stale delivery still succeeds at the object level
//! (a failure would loop through MRF re-delivering the stale value).
use crate::common::{RustFSTestEnvironment, init_logging};
use aws_sdk_s3::Client;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use serial_test::serial;
type TestResult = Result<(), Box<dyn std::error::Error + Send + Sync>>;
const HDR_SOURCE_REPLICATION_REQUEST: &str = "x-rustfs-source-replication-request";
const HDR_SOURCE_VERSION_ID: &str = "x-rustfs-source-version-id";
const HDR_SOURCE_MTIME: &str = "x-rustfs-source-mtime";
const HDR_SOURCE_TAGGING_TIMESTAMP: &str = "x-rustfs-source-replication-tagging-timestamp";
const SOURCE_MTIME: &str = "2026-01-01T00:00:00Z";
const T_STALE: &str = "2026-01-01T00:00:01Z";
const T_LOCAL: &str = "2026-02-01T00:00:00Z";
const T_NEWER: &str = "2026-03-01T00:00:00Z";
/// Simulated inbound authorized replication PUT: same object version, tags and
/// the source-authored tagging timestamp carried in transport headers.
async fn inbound_replication_put(
client: &Client,
bucket: &str,
key: &str,
version_id: &str,
tags: &str,
tagging_timestamp: &str,
) -> TestResult {
let version_id = version_id.to_string();
let tagging_timestamp = tagging_timestamp.to_string();
client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from_static(b"lww-e2e-body"))
.tagging(tags)
.customize()
.mutate_request(move |req| {
req.headers_mut().insert(HDR_SOURCE_REPLICATION_REQUEST, "true");
req.headers_mut().insert(HDR_SOURCE_VERSION_ID, version_id.clone());
req.headers_mut().insert(HDR_SOURCE_MTIME, SOURCE_MTIME);
req.headers_mut()
.insert(HDR_SOURCE_TAGGING_TIMESTAMP, tagging_timestamp.clone());
})
.send()
.await?;
Ok(())
}
async fn tag_value(client: &Client, bucket: &str, key: &str, version_id: &str, tag_key: &str) -> Option<String> {
let tagging = client
.get_object_tagging()
.bucket(bucket)
.key(key)
.version_id(version_id)
.send()
.await
.expect("object tagging should be readable");
tagging
.tag_set()
.iter()
.find(|tag| tag.key() == tag_key)
.map(|tag| tag.value().to_string())
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn receiver_lww_keeps_newer_tags_across_delivery_orders() -> TestResult {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let client = env.create_s3_client();
let bucket = "replication-lww-receiver";
let key = "object";
client.create_bucket().bucket(bucket).send().await?;
client
.put_bucket_versioning()
.bucket(bucket)
.versioning_configuration(
VersioningConfiguration::builder()
.status(BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await?;
// First delivery establishes version V with tags stamped T_LOCAL.
let version_id = uuid::Uuid::new_v4().to_string();
inbound_replication_put(&client, bucket, key, &version_id, "site=local", T_LOCAL).await?;
assert_eq!(
tag_value(&client, bucket, key, &version_id, "site").await.as_deref(),
Some("local"),
"the first delivery must establish the tagged version"
);
// A stale delivery (older source timestamp) must succeed at the object
// level but must NOT overwrite the newer tags.
inbound_replication_put(&client, bucket, key, &version_id, "site=stale", T_STALE).await?;
assert_eq!(
tag_value(&client, bucket, key, &version_id, "site").await.as_deref(),
Some("local"),
"a stale inbound delivery must not overwrite newer tags (rustfs/backlog#1953)"
);
// A newer delivery still converges the version onto the newest state.
inbound_replication_put(&client, bucket, key, &version_id, "site=newer", T_NEWER).await?;
assert_eq!(
tag_value(&client, bucket, key, &version_id, "site").await.as_deref(),
Some("newer"),
"a newer inbound delivery must overwrite older tags"
);
client
.delete_object()
.bucket(bucket)
.key(key)
.version_id(&version_id)
.send()
.await?;
env.delete_test_bucket(bucket).await.ok();
Ok(())
}
+6 -6
View File
@@ -199,12 +199,12 @@ pub mod bucket {
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, get_proxy_targets, init_background_replication,
invalid_replication_config_status_field, is_site_replication_rule, merge_incoming_replication_config,
persist_force_delete_intent, read_durable_mrf_backlog, replication_state_to_filemeta, replication_status_to_filemeta,
replication_statuses_map, replication_target_arn_deployment_id, replication_target_arns, resync_start_conflict_id,
should_remove_replication_target, should_schedule_delete_replication, should_use_existing_delete_replication_info,
should_use_existing_delete_replication_source, unsupported_replication_config_field,
validate_replication_config_structure, validate_replication_config_target_arns, version_purge_status_to_filemeta,
invalid_replication_config_status_field, persist_force_delete_intent, read_durable_mrf_backlog,
replication_state_to_filemeta, replication_status_to_filemeta, replication_statuses_map, replication_target_arns,
resync_start_conflict_id, should_remove_replication_target, should_schedule_delete_replication,
should_use_existing_delete_replication_info, should_use_existing_delete_replication_source,
unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns,
version_purge_status_to_filemeta,
};
}
+1 -2
View File
@@ -47,8 +47,7 @@ pub use replication_config_boundary::{
ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS,
REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS,
ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationTargetValidationError,
invalid_replication_config_status_field, is_site_replication_rule, merge_incoming_replication_config,
replication_target_arn_deployment_id, replication_target_arns, should_remove_replication_target,
invalid_replication_config_status_field, replication_target_arns, should_remove_replication_target,
unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns,
};
pub(crate) use replication_filemeta_boundary::version_purge_statuses_map;
@@ -16,7 +16,6 @@ pub use rustfs_replication::{
ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS,
REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS,
ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationRuleExt, ReplicationTargetValidationError,
invalid_replication_config_status_field, is_site_replication_rule, merge_incoming_replication_config,
replication_target_arn_deployment_id, replication_target_arns, should_remove_replication_target,
invalid_replication_config_status_field, replication_target_arns, should_remove_replication_target,
unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns,
};
@@ -3868,6 +3868,7 @@ async fn replicate_object_with_multipart<S: ReplicationObjectIO>(ctx: MultipartR
actual_size,
object_info.etag.clone().unwrap_or_default(),
object_info.mod_time,
&put_opts.internal,
),
)
.await
@@ -472,6 +472,7 @@ pub(crate) fn replication_complete_multipart_options(
actual_size: String,
source_etag: String,
source_mtime: Option<OffsetDateTime>,
source_internal: &AdvancedPutOptions,
) -> PutObjectOptions {
let mut user_metadata = HashMap::new();
insert_header_map(&mut user_metadata, SUFFIX_REPLICATION_ACTUAL_OBJECT_SIZE, actual_size);
@@ -484,6 +485,14 @@ pub(crate) fn replication_complete_multipart_options(
// mtime must degrade to epoch so header() suppresses the header
// instead of asserting the replication time as the object's mtime.
source_mtime: source_mtime.unwrap_or(OffsetDateTime::UNIX_EPOCH),
// Carry the per-category LWW timestamps on the complete request as
// well: the receiver's CompleteMultipartUpload options builder
// parses the same headers, so the multipart transport gets the
// same receiver-side LWW as the single-PUT transport
// (rustfs/backlog#1953). Epoch values keep the headers suppressed.
tagging_timestamp: source_internal.tagging_timestamp,
retention_timestamp: source_internal.retention_timestamp,
legalhold_timestamp: source_internal.legalhold_timestamp,
replication_status: ReplicationStatusType::Replica,
replication_request: true,
..Default::default()
@@ -649,20 +658,39 @@ mod tests {
#[test]
fn replication_complete_multipart_options_sets_actual_size() {
let source_mtime = OffsetDateTime::from_unix_timestamp(1_716_170_000).expect("valid test timestamp");
let source_internal = AdvancedPutOptions {
tagging_timestamp: OffsetDateTime::from_unix_timestamp(1_716_170_100).expect("valid test timestamp"),
retention_timestamp: OffsetDateTime::from_unix_timestamp(1_716_170_200).expect("valid test timestamp"),
legalhold_timestamp: OffsetDateTime::from_unix_timestamp(1_716_170_300).expect("valid test timestamp"),
..Default::default()
};
let options = replication_complete_multipart_options(
"1024".to_string(),
"0123456789abcdef0123456789abcdef-3".to_string(),
Some(source_mtime),
&source_internal,
);
assert_eq!(options.internal.source_etag, "0123456789abcdef0123456789abcdef-3");
assert_eq!(options.internal.source_mtime, source_mtime);
// The complete request must carry the same per-category LWW timestamps
// as the initiate request; the receiver reads them from the complete
// headers (rustfs/backlog#1953).
assert_eq!(options.internal.tagging_timestamp, source_internal.tagging_timestamp);
assert_eq!(options.internal.retention_timestamp, source_internal.retention_timestamp);
assert_eq!(options.internal.legalhold_timestamp, source_internal.legalhold_timestamp);
// Absent source mtime must degrade to epoch (header suppressed), not
// the AdvancedPutOptions default of now_utc() — that default would
// stamp the replication time as the replica's mtime and break the
// multipart HEAD convergence.
let options_no_mtime = replication_complete_multipart_options("1024".to_string(), String::new(), None);
// multipart HEAD convergence. Unset category timestamps stay epoch so
// header() keeps suppressing them.
let options_no_mtime =
replication_complete_multipart_options("1024".to_string(), String::new(), None, &AdvancedPutOptions::default());
assert_eq!(options_no_mtime.internal.source_mtime.unix_timestamp(), 0);
assert_eq!(options_no_mtime.internal.tagging_timestamp.unix_timestamp(), 0);
assert_eq!(options_no_mtime.internal.retention_timestamp.unix_timestamp(), 0);
assert_eq!(options_no_mtime.internal.legalhold_timestamp.unix_timestamp(), 0);
assert_eq!(
get_header_map(&options.user_metadata, SUFFIX_REPLICATION_ACTUAL_OBJECT_SIZE).as_deref(),
@@ -2223,6 +2223,57 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
fi.set_data_moved();
}
// Receiver-side LWW (rustfs/backlog#1953): the multipart replication
// transport carries the category values at CreateMultipartUpload (in
// the staged upload metadata) and the source category timestamps on
// the complete request. Read the destination version under the held
// object write lock and keep any category this site modified more
// recently. Read failures (version absent on first replication, quorum
// errors) keep today's overwrite semantics: failing the complete would
// loop through MRF, re-delivering the stale value forever.
if crate::set_disk::ops::object::replication_lww_applicable(opts)
&& let Some(version_id) = fi.version_id
{
match self
.get_object_info(
bucket,
object,
&ObjectOptions {
version_id: Some(version_id.to_string()),
no_lock: true,
metadata_cache_safe: false,
versioned: opts.versioned,
version_suspended: opts.version_suspended,
..Default::default()
},
)
.await
{
Ok(existing) => {
let stored = crate::set_disk::ops::object::stored_replication_category_metadata(&existing);
crate::set_disk::ops::object::merge_replication_metadata_lww(&mut fi.metadata, &stored, opts);
}
// Version absent: first replication of this version, nothing
// local to compare — the normal path, not a degraded one.
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {}
Err(err) => {
// Degraded path: without the stored state the inbound
// metadata is applied unchanged — exactly the overwrite
// LWW exists to prevent — so this must be operator-visible.
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
version_id = %version_id,
error = %err,
state = "replication_lww_read_unavailable",
"SetDisk multipart replication LWW read skipped; inbound metadata applied without comparison"
);
}
}
}
for meta in parts_metadatas.iter_mut() {
if meta.has_valid_erasure_geometry() {
meta.size = fi.size;
@@ -6960,6 +7011,97 @@ mod tests {
.await
}
/// Receiver-side LWW on the multipart replication transport
/// (rustfs/backlog#1953): a metadata-only replication of a multipart
/// source object rides CreateMultipartUpload (category values in the
/// upload metadata) + CompleteMultipartUpload (category timestamps in
/// the complete options). A stale inbound tagging timestamp must not
/// overwrite a newer locally-tagged destination version.
#[tokio::test]
#[serial]
async fn complete_multipart_upload_stale_replication_tags_keep_local() {
use rustfs_utils::http::headers::AMZ_OBJECT_TAGGING;
use rustfs_utils::http::{SUFFIX_TAGGING_TIMESTAMP, get_str};
use time::format_description::well_known::Rfc3339;
const T_OLD: &str = "2026-01-01T00:00:00Z";
const T_LOCAL: &str = "2026-02-01T00:00:00Z";
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "multipart-replication-lww-bucket";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
// Local destination version with newer tags.
let version_id = Uuid::new_v4();
let mut local_metadata = HashMap::new();
local_metadata.insert(AMZ_OBJECT_TAGGING.to_string(), "site=local".to_string());
rustfs_utils::http::insert_str(&mut local_metadata, SUFFIX_TAGGING_TIMESTAMP, T_LOCAL.to_string());
let mut local_reader = PutObjReader::from_vec(b"local body".to_vec());
set_disks
.put_object(
bucket,
object,
&mut local_reader,
&ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
user_defined: local_metadata,
// Explicit-version PUTs require the bucket Object Lock snapshot.
object_lock_config_snapshot: Some(Arc::new(crate::set_disk::ObjectLockConfigSnapshot::new(
crate::bucket::metadata_sys::ObjectLockConfigState::ConfirmedAbsent,
))),
..Default::default()
},
)
.await
.expect("local versioned put should commit");
// Inbound replication upload carrying older tags for the same version.
let mut inbound_metadata = HashMap::new();
inbound_metadata.insert(AMZ_OBJECT_TAGGING.to_string(), "site=remote".to_string());
rustfs_utils::http::insert_str(&mut inbound_metadata, SUFFIX_TAGGING_TIMESTAMP, T_OLD.to_string());
let create_opts = ObjectOptions {
versioned: true,
user_defined: inbound_metadata,
..Default::default()
};
let (upload_id, parts) =
stage_upload_with_create_opts(&set_disks, bucket, object, &payload(0x5a), &create_opts).await;
rewrite_staged_upload_version_id(&set_disks, bucket, object, &upload_id, Some(version_id)).await;
let complete_opts = ObjectOptions {
versioned: true,
replication_request: true,
replication_tagging_timestamp: Some(OffsetDateTime::parse(T_OLD, &Rfc3339).expect("test timestamp should parse")),
..Default::default()
};
set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, parts, &complete_opts)
.await
.expect("replication multipart completion should succeed even when a category keeps local values");
let info = set_disks
.get_object_info(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
..Default::default()
},
)
.await
.expect("completed version should be readable");
assert_eq!(
info.user_tags.as_str(),
"site=local",
"older inbound multipart tags must not overwrite newer local tags"
);
assert_eq!(get_str(&info.user_defined, SUFFIX_TAGGING_TIMESTAMP).as_deref(), Some(T_LOCAL));
}
#[tokio::test]
#[serial]
async fn complete_multipart_upload_assigns_completion_version_id() {
+471
View File
@@ -1880,6 +1880,110 @@ fn delete_file_info_with_replication_transport_metadata(fi: &FileInfo) -> FileIn
transported
}
/// True when an authorized replication write carries at least one per-category
/// source timestamp, i.e. receiver-side LWW has something to judge.
pub(in crate::set_disk) fn replication_lww_applicable(opts: &ObjectOptions) -> bool {
opts.replication_request
&& (opts.replication_tagging_timestamp.is_some()
|| opts.replication_retention_timestamp.is_some()
|| opts.replication_legalhold_timestamp.is_some())
}
/// The stored per-category state of a destination version, as compared by
/// [`merge_replication_metadata_lww`]. `ObjectInfo::from_file_info`
/// externalizes tags into `user_tags` (stripping the metadata key), so the
/// tag value is folded back into map form here.
pub(in crate::set_disk) fn stored_replication_category_metadata(existing: &ObjectInfo) -> HashMap<String, String> {
let mut stored = (*existing.user_defined).clone();
if !existing.user_tags.is_empty() {
stored.insert(rustfs_utils::http::headers::AMZ_OBJECT_TAGGING.to_string(), (*existing.user_tags).clone());
}
stored
}
/// Receiver-side last-writer-wins for authorized replication writes
/// (rustfs/backlog#1953, audit A4/P1-6). Metadata-only replication reuses the
/// whole-object transports, so in active-active topologies an inbound write
/// carries the source's tags / retention / legal hold verbatim and would
/// otherwise overwrite a category the destination modified more recently —
/// both sites end up permanently diverged while reporting COMPLETED.
///
/// Judged per category, only when the inbound request carries that category's
/// source timestamp (`ObjectOptions::replication_*_timestamp`):
/// - stored timestamp newer than inbound: the local category values and
/// timestamp are kept; the rest of the write proceeds per the inbound
/// metadata and the object-level result stays successful (failing the write
/// instead would loop through MRF, re-delivering the stale value forever);
/// - otherwise the inbound category wins and its internal timestamp key is
/// pinned to the source-authored time — the PUT path re-stamps the
/// object-lock timestamps with the receiver's clock
/// (`parse_object_lock_retention` / `parse_object_lock_legal_hold` insert
/// `now()` via `eval_metadata`), which would make the replica's clock the
/// LWW authority and wedge later convergence;
/// - no stored timestamp (pre-P1-6 data) or no inbound timestamp: the current
/// overwrite behavior is preserved.
///
/// Returns whether `inbound` was modified. Callers must hold the object write
/// lock so the stored values compared here are the ones being replaced.
pub(in crate::set_disk) fn merge_replication_metadata_lww(
inbound: &mut HashMap<String, String>,
existing: &HashMap<String, String>,
opts: &ObjectOptions,
) -> bool {
use rustfs_utils::http::headers::{
AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER, AMZ_OBJECT_LOCK_MODE_LOWER, AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER, AMZ_OBJECT_TAGGING,
};
use rustfs_utils::http::metadata_compat::{
SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, SUFFIX_TAGGING_TIMESTAMP, get_str,
remove_str,
};
use time::format_description::well_known::Rfc3339;
let categories: [(Option<OffsetDateTime>, &str, &[&str]); 3] = [
(opts.replication_tagging_timestamp, SUFFIX_TAGGING_TIMESTAMP, &[AMZ_OBJECT_TAGGING]),
(
opts.replication_retention_timestamp,
SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP,
&[AMZ_OBJECT_LOCK_MODE_LOWER, AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER],
),
(
opts.replication_legalhold_timestamp,
SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP,
&[AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER],
),
];
let mut changed = false;
for (inbound_timestamp, timestamp_suffix, value_keys) in categories {
let Some(inbound_timestamp) = inbound_timestamp else { continue };
let is_category_value_key = |key: &str| value_keys.iter().any(|value_key| key.eq_ignore_ascii_case(value_key));
let stored_timestamp = get_str(existing, timestamp_suffix).and_then(|value| OffsetDateTime::parse(&value, &Rfc3339).ok());
if stored_timestamp.is_some_and(|stored| stored > inbound_timestamp) {
inbound.retain(|key, _| !is_category_value_key(key));
remove_str(inbound, timestamp_suffix);
for (key, value) in existing {
if is_category_value_key(key) {
inbound.insert(key.clone(), value.clone());
}
}
// Restore the winning timestamp via insert_str, not a verbatim key
// copy: a MinIO-written version may carry only the
// x-minio-internal- key, and the dual-key invariant requires every
// write to produce both keys.
if let Some(stored_value) = get_str(existing, timestamp_suffix) {
rustfs_utils::http::insert_str(inbound, timestamp_suffix, stored_value);
}
changed = true;
} else if let Ok(source_authored) = inbound_timestamp.format(&Rfc3339)
&& get_str(inbound, timestamp_suffix).as_deref() != Some(source_authored.as_str())
{
rustfs_utils::http::insert_str(inbound, timestamp_suffix, source_authored);
changed = true;
}
}
changed
}
impl SetDisks {
pub(in crate::set_disk) async fn persist_old_data_cleanup_receipts(
&self,
@@ -2560,6 +2664,22 @@ impl SetDisks {
if check_object_lock_for_deletion_with_state(object_lock_config.state(), &existing, false)?.is_some() {
return Err(StorageError::PrefixAccessDenied(bucket.to_string(), object.to_string()));
}
// Receiver-side LWW (rustfs/backlog#1953): reuse this
// commit-lock read of the destination version so a
// category (tags / retention / legal hold) modified
// more recently on this site is kept instead of being
// overwritten by the inbound replication metadata.
if replication_lww_applicable(opts) {
let stored = stored_replication_category_metadata(&existing);
let mut merged = parts_metadatas[response_metadata_slot].metadata.clone();
if merge_replication_metadata_lww(&mut merged, &stored, opts) {
for (pfi, disk) in parts_metadatas.iter_mut().zip(shuffle_disks.iter()) {
if disk.is_some() {
pfi.metadata = merged.clone();
}
}
}
}
}
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {}
Err(err) => return Err(err),
@@ -7887,6 +8007,357 @@ mod replication_quota_safety_tests {
}
}
#[cfg(test)]
mod replication_lww_tests {
//! Receiver-side LWW for authorized replication writes (rustfs/backlog#1953,
//! audit A4/P1-6): an inbound replication PUT whose per-category timestamp
//! (tags / retention / legal hold) is older than the destination version's
//! stored timestamp must keep the local category values instead of
//! overwriting them; categories are judged independently and the write
//! itself still succeeds.
use super::hermetic_set_disks_support::hermetic_set_disks_isolated as hermetic_set_disks;
use super::*;
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
use rustfs_utils::http::headers::{
AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER, AMZ_OBJECT_LOCK_MODE_LOWER, AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER, AMZ_OBJECT_TAGGING,
};
use rustfs_utils::http::{
SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, SUFFIX_TAGGING_TIMESTAMP, get_str,
insert_str,
};
use time::format_description::well_known::Rfc3339;
const T_OLD: &str = "2026-01-01T00:00:00Z";
const T_LOCAL: &str = "2026-02-01T00:00:00Z";
const T_NEW: &str = "2026-03-01T00:00:00Z";
fn parse_ts(value: &str) -> OffsetDateTime {
OffsetDateTime::parse(value, &Rfc3339).expect("test timestamp should parse")
}
async fn make_bucket(disks: &[DiskStore], bucket: &str) {
for disk in disks {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
}
async fn put_version(set_disks: &Arc<SetDisks>, bucket: &str, object: &str, version_id: &str, opts: &ObjectOptions) {
let mut reader = PutObjReader::from_vec(b"lww-body".to_vec());
set_disks
.put_object(bucket, object, &mut reader, opts)
.await
.expect("versioned put should commit");
assert_eq!(opts.version_id.as_deref(), Some(version_id));
}
fn versioned_opts(version_id: &str, user_defined: HashMap<String, String>) -> ObjectOptions {
ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
user_defined,
// Explicit-version PUTs require the bucket Object Lock snapshot.
object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(
crate::bucket::metadata_sys::ObjectLockConfigState::ConfirmedAbsent,
))),
..Default::default()
}
}
/// Local state: version `version_id` with tags "site=local" stamped `T_LOCAL`.
async fn seed_local_tagged_version(set_disks: &Arc<SetDisks>, bucket: &str, object: &str, version_id: &str) {
let mut user_defined = HashMap::new();
user_defined.insert(AMZ_OBJECT_TAGGING.to_string(), "site=local".to_string());
insert_str(&mut user_defined, SUFFIX_TAGGING_TIMESTAMP, T_LOCAL.to_string());
put_version(set_disks, bucket, object, version_id, &versioned_opts(version_id, user_defined)).await;
}
fn inbound_tagging_opts(version_id: &str, tags: &str, timestamp: &str) -> ObjectOptions {
let mut user_defined = HashMap::new();
user_defined.insert(AMZ_OBJECT_TAGGING.to_string(), tags.to_string());
insert_str(&mut user_defined, SUFFIX_TAGGING_TIMESTAMP, timestamp.to_string());
ObjectOptions {
replication_request: true,
replication_tagging_timestamp: Some(parse_ts(timestamp)),
..versioned_opts(version_id, user_defined)
}
}
async fn version_info(set_disks: &Arc<SetDisks>, bucket: &str, object: &str, version_id: &str) -> ObjectInfo {
set_disks
.get_object_info(bucket, object, &versioned_opts(version_id, HashMap::new()))
.await
.expect("version should be readable")
}
#[tokio::test]
async fn inbound_stale_tagging_keeps_newer_local_tags() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-tagging-stale";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
seed_local_tagged_version(&set_disks, bucket, object, &version_id).await;
put_version(
&set_disks,
bucket,
object,
&version_id,
&inbound_tagging_opts(&version_id, "site=remote", T_OLD),
)
.await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert_eq!(
info.user_tags.as_str(),
"site=local",
"older inbound tags must not overwrite newer local tags"
);
assert_eq!(
get_str(&info.user_defined, SUFFIX_TAGGING_TIMESTAMP).as_deref(),
Some(T_LOCAL),
"the winning local tagging timestamp must be preserved"
);
}
#[tokio::test]
async fn inbound_newer_tagging_overwrites_local_tags() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-tagging-newer";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
seed_local_tagged_version(&set_disks, bucket, object, &version_id).await;
put_version(
&set_disks,
bucket,
object,
&version_id,
&inbound_tagging_opts(&version_id, "site=remote", T_NEW),
)
.await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert_eq!(
info.user_tags.as_str(),
"site=remote",
"newer inbound tags must overwrite older local tags"
);
assert_eq!(get_str(&info.user_defined, SUFFIX_TAGGING_TIMESTAMP).as_deref(), Some(T_NEW));
}
#[tokio::test]
async fn inbound_wins_when_local_has_no_tagging_timestamp() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-tagging-no-local-ts";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
// Pre-P1-6 data: local tags without a stored tagging timestamp.
let mut user_defined = HashMap::new();
user_defined.insert(AMZ_OBJECT_TAGGING.to_string(), "site=local".to_string());
put_version(&set_disks, bucket, object, &version_id, &versioned_opts(&version_id, user_defined)).await;
put_version(
&set_disks,
bucket,
object,
&version_id,
&inbound_tagging_opts(&version_id, "site=remote", T_OLD),
)
.await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert_eq!(
info.user_tags.as_str(),
"site=remote",
"without a local timestamp the inbound category must win (pre-LWW data compatibility)"
);
}
#[tokio::test]
async fn categories_are_judged_independently() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-category-independent";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
// Local: newer tags (T_LOCAL), older *cleared* retention (T_OLD) —
// timestamp key only, the shape a replicated retention clear stores.
// (An active local retention would already block the overwrite at the
// WORM gate; the LWW-reachable retention states are cleared/expired.)
let mut local = HashMap::new();
local.insert(AMZ_OBJECT_TAGGING.to_string(), "site=local".to_string());
insert_str(&mut local, SUFFIX_TAGGING_TIMESTAMP, T_LOCAL.to_string());
insert_str(&mut local, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, T_OLD.to_string());
put_version(&set_disks, bucket, object, &version_id, &versioned_opts(&version_id, local)).await;
// Inbound: older tags (T_OLD), newer retention (T_NEW).
let mut inbound = HashMap::new();
inbound.insert(AMZ_OBJECT_TAGGING.to_string(), "site=remote".to_string());
insert_str(&mut inbound, SUFFIX_TAGGING_TIMESTAMP, T_OLD.to_string());
inbound.insert(AMZ_OBJECT_LOCK_MODE_LOWER.to_string(), "COMPLIANCE".to_string());
inbound.insert(AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER.to_string(), "2028-01-01T00:00:00Z".to_string());
insert_str(&mut inbound, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, T_NEW.to_string());
let opts = ObjectOptions {
replication_request: true,
replication_tagging_timestamp: Some(parse_ts(T_OLD)),
replication_retention_timestamp: Some(parse_ts(T_NEW)),
..versioned_opts(&version_id, inbound)
};
put_version(&set_disks, bucket, object, &version_id, &opts).await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert_eq!(info.user_tags.as_str(), "site=local", "the stale tagging category must keep local values");
assert_eq!(
info.user_defined.get(AMZ_OBJECT_LOCK_MODE_LOWER).map(String::as_str),
Some("COMPLIANCE"),
"the newer retention category must be applied in the same write"
);
assert_eq!(get_str(&info.user_defined, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP).as_deref(), Some(T_NEW));
}
#[tokio::test]
async fn inbound_stale_legal_hold_keeps_local_value() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-legalhold-stale";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
// Local: legal hold released (OFF) at T_LOCAL. (A local hold that is
// still ON already blocks the overwrite at the WORM gate; the
// LWW-reachable divergence is a stale inbound ON resurrecting a hold
// that was released more recently on this site.)
let mut local = HashMap::new();
local.insert(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER.to_string(), "OFF".to_string());
insert_str(&mut local, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP, T_LOCAL.to_string());
put_version(&set_disks, bucket, object, &version_id, &versioned_opts(&version_id, local)).await;
let mut inbound = HashMap::new();
inbound.insert(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER.to_string(), "ON".to_string());
insert_str(&mut inbound, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP, T_OLD.to_string());
let opts = ObjectOptions {
replication_request: true,
replication_legalhold_timestamp: Some(parse_ts(T_OLD)),
..versioned_opts(&version_id, inbound)
};
put_version(&set_disks, bucket, object, &version_id, &opts).await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert_eq!(
info.user_defined.get(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER).map(String::as_str),
Some("OFF"),
"a stale inbound legal hold must not resurrect a hold released more recently"
);
assert_eq!(
get_str(&info.user_defined, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP).as_deref(),
Some(T_LOCAL)
);
}
/// Dual-key invariant under LWW: a MinIO-written destination version may
/// carry only the x-minio-internal timestamp key; when the local category
/// wins, the restored map must still hold BOTH compatibility keys.
#[test]
fn local_win_restores_both_internal_timestamp_keys_for_minio_only_metadata() {
let mut inbound = HashMap::new();
inbound.insert(AMZ_OBJECT_TAGGING.to_string(), "site=remote".to_string());
insert_str(&mut inbound, SUFFIX_TAGGING_TIMESTAMP, T_OLD.to_string());
let existing = HashMap::from([
(AMZ_OBJECT_TAGGING.to_string(), "site=local".to_string()),
("X-Minio-Internal-Tagging-Timestamp".to_string(), T_LOCAL.to_string()),
]);
let opts = ObjectOptions {
replication_request: true,
replication_tagging_timestamp: Some(parse_ts(T_OLD)),
..Default::default()
};
assert!(merge_replication_metadata_lww(&mut inbound, &existing, &opts));
assert_eq!(inbound.get(AMZ_OBJECT_TAGGING).map(String::as_str), Some("site=local"));
assert_eq!(
inbound.get("x-rustfs-internal-tagging-timestamp").map(String::as_str),
Some(T_LOCAL),
"the RustFS twin key must be materialized even when the source version only had the MinIO key"
);
assert_eq!(inbound.get("x-minio-internal-tagging-timestamp").map(String::as_str), Some(T_LOCAL));
}
/// When the inbound category wins, the stored timestamp must be the
/// source-authored one: the PUT path's eval_metadata stamps the
/// object-lock timestamps with the receiver's clock
/// (`parse_object_lock_retention`), which would otherwise make this
/// replica's clock the LWW authority and wedge later convergence.
#[tokio::test]
async fn inbound_win_pins_stored_timestamp_to_source_authored_value() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-retention-ts-pinned";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
// Local cleared retention at T_OLD.
let mut local = HashMap::new();
insert_str(&mut local, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, T_OLD.to_string());
put_version(&set_disks, bucket, object, &version_id, &versioned_opts(&version_id, local)).await;
// Inbound newer retention: the source authored T_LOCAL, but the PUT
// path's eval_metadata stomped the metadata key with receiver-now
// (simulated by T_NEW here).
let mut inbound = HashMap::new();
inbound.insert(AMZ_OBJECT_LOCK_MODE_LOWER.to_string(), "GOVERNANCE".to_string());
inbound.insert(AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER.to_string(), "2028-01-01T00:00:00Z".to_string());
insert_str(&mut inbound, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, T_NEW.to_string());
let opts = ObjectOptions {
replication_request: true,
replication_retention_timestamp: Some(parse_ts(T_LOCAL)),
..versioned_opts(&version_id, inbound)
};
put_version(&set_disks, bucket, object, &version_id, &opts).await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert_eq!(
get_str(&info.user_defined, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP).as_deref(),
Some(T_LOCAL),
"the stored category timestamp must be the source-authored time, not the receiver's clock"
);
assert_eq!(info.user_defined.get(AMZ_OBJECT_LOCK_MODE_LOWER).map(String::as_str), Some("GOVERNANCE"));
}
#[tokio::test]
async fn newer_local_tag_deletion_survives_stale_inbound_tags() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "lww-tagging-deleted";
let object = "object";
let version_id = Uuid::new_v4().to_string();
make_bucket(&disk_stores, bucket).await;
// Local DeleteObjectTagging state: no tags, but a newer tagging timestamp.
let mut local = HashMap::new();
insert_str(&mut local, SUFFIX_TAGGING_TIMESTAMP, T_LOCAL.to_string());
put_version(&set_disks, bucket, object, &version_id, &versioned_opts(&version_id, local)).await;
put_version(
&set_disks,
bucket,
object,
&version_id,
&inbound_tagging_opts(&version_id, "site=remote", T_OLD),
)
.await;
let info = version_info(&set_disks, bucket, object, &version_id).await;
assert!(
info.user_tags.is_empty(),
"a newer local tag deletion must not be resurrected by older inbound tags"
);
assert_eq!(get_str(&info.user_defined, SUFFIX_TAGGING_TIMESTAMP).as_deref(), Some(T_LOCAL));
}
}
#[cfg(test)]
mod inline_put_commit_path_tests {
use super::hermetic_set_disks_support::hermetic_set_disks_isolated as hermetic_set_disks;
-72
View File
@@ -265,78 +265,6 @@ pub fn active_replication_rule_destination_arns(config: &ReplicationConfiguratio
arns
}
/// Deployment id extracted from a site-replication target ARN
/// (`arn:{rustfs|minio}:replication::<deployment-id>:<bucket>`), or `None`
/// for an operator-authored ARN.
pub fn replication_target_arn_deployment_id(arn: &str) -> Option<String> {
let parts: Vec<_> = arn.split(':').collect();
if parts.len() == 6
&& parts[0] == "arn"
&& matches!(parts[1], "rustfs" | "minio")
&& parts[2] == "replication"
&& !parts[4].is_empty()
{
return Some(parts[4].to_string());
}
None
}
/// Whether `rule` is a site-replication rule (`site-repl-*` id) owned by the
/// local site's reconciler rather than authored by an operator.
pub fn is_site_replication_rule(rule: &ReplicationRule) -> bool {
rule.id.as_deref().is_some_and(|id| id.starts_with("site-repl-"))
}
/// Merge an incoming replication config into the local one.
///
/// `site-repl-*` rules encode the *holder's* outbound direction — their
/// destination ARN names another site — so applying an external rule set
/// verbatim replaces the local reverse rule with one this site can never
/// satisfy (no bucket target backs it) and replication silently stops. Only
/// operator-authored rules travel: the site-replication peer ingestion path
/// and the S3 put/delete-bucket-replication path both keep the local site's
/// `site-repl-*` rules through this merge. `incoming == None` models a
/// delete of the operator-authored rules.
pub fn merge_incoming_replication_config(
incoming: Option<ReplicationConfiguration>,
local: Option<ReplicationConfiguration>,
) -> Option<ReplicationConfiguration> {
let incoming_role = incoming.as_ref().map(|config| config.role.clone()).unwrap_or_default();
// Operator rules first, then the local site rules — the same order the
// site-replication reconciler produces, so its no-op check matches and
// the bucket metadata is written once per broadcast, not twice.
let mut rules: Vec<ReplicationRule> = incoming
.into_iter()
.flat_map(|config| config.rules)
.filter(|rule| !is_site_replication_rule(rule))
.collect();
rules.extend(
local
.into_iter()
.flat_map(|config| config.rules)
.filter(is_site_replication_rule),
);
if rules.is_empty() {
return None;
}
for (index, rule) in rules.iter_mut().enumerate() {
rule.priority = Some(i32::try_from(index + 1).unwrap_or(i32::MAX));
}
// A site-replication ARN in `role` is the sender's, and the reconciler's
// per-peer target lookup reads it — carrying it over would pin the
// receiver's targets to the sender's identity.
let role = match replication_target_arn_deployment_id(&incoming_role) {
Some(_) => String::new(),
None => incoming_role,
};
Some(ReplicationConfiguration { role, rules })
}
pub fn replication_target_arns(config: &ReplicationConfiguration) -> HashSet<String> {
let role = config.role.trim();
if !role.is_empty() {
+1 -2
View File
@@ -32,8 +32,7 @@ pub use config::{
ObjectOpts, REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS,
REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS,
ReplicationConfigStructureError, ReplicationConfigurationExt, ReplicationTargetValidationError,
active_replication_rule_destination_arns, invalid_replication_config_status_field, is_site_replication_rule,
merge_incoming_replication_config, replication_target_arn_deployment_id, replication_target_arns,
active_replication_rule_destination_arns, invalid_replication_config_status_field, replication_target_arns,
should_remove_replication_target, unsupported_replication_config_field, validate_replication_config_structure,
validate_replication_config_target_arns,
};
+64 -11
View File
@@ -31,9 +31,6 @@ use crate::admin::storage_api::bucket::metadata::{
use crate::admin::storage_api::bucket::metadata_sys;
use crate::admin::storage_api::bucket::quota::BucketQuota;
use crate::admin::storage_api::bucket::replication;
use crate::admin::storage_api::bucket::replication::{
is_site_replication_rule, merge_incoming_replication_config, replication_target_arn_deployment_id,
};
use crate::admin::storage_api::bucket::target::{ARN, BucketTarget, BucketTargetType, BucketTargets, Credentials};
use crate::admin::storage_api::bucket::target_sys::BucketTargetSys;
use crate::admin::storage_api::bucket::utils::{deserialize, serialize};
@@ -1120,14 +1117,6 @@ async fn load_site_replication_state() -> S3Result<SiteReplicationState> {
}
}
/// Whether this deployment participates in site replication (two or more
/// peers in the persisted state). Read by the S3 interface layer to gate
/// replication-config edits (MinIO `ErrReplicationDenyEditError` semantics,
/// issue #1948); a state-read failure propagates so the gate fails closed.
pub(crate) async fn site_replication_enabled() -> S3Result<bool> {
Ok(load_site_replication_state().await?.enabled())
}
async fn load_site_replication_state_no_lock(store: Arc<ECStore>) -> S3Result<SiteReplicationState> {
match read_config_no_lock(store, SITE_REPLICATION_STATE_PATH).await {
Ok(data) => parse_site_replication_state(&data),
@@ -7759,6 +7748,20 @@ fn bucket_target_deployment_id(target: &BucketTarget) -> Option<String> {
replication_target_arn_deployment_id(&target.arn)
}
fn replication_target_arn_deployment_id(arn: &str) -> Option<String> {
let parts: Vec<_> = arn.split(':').collect();
if parts.len() == 6
&& parts[0] == "arn"
&& matches!(parts[1], "rustfs" | "minio")
&& parts[2] == "replication"
&& !parts[4].is_empty()
{
return Some(parts[4].to_string());
}
None
}
fn prune_removed_site_replication_bucket_targets(
existing: BucketTargets,
removed_deployment_ids: &HashSet<String>,
@@ -7783,6 +7786,10 @@ fn prune_removed_site_replication_bucket_targets(
(BucketTargets { targets }, removed)
}
fn is_site_replication_rule(rule: &ReplicationRule) -> bool {
rule.id.as_deref().is_some_and(|id| id.starts_with("site-repl-"))
}
/// Whether every `site-repl-*` rule on this bucket resolves to a live remote target.
///
/// The rule set alone cannot answer this: a rule can be perfectly formed while the endpoint
@@ -7808,6 +7815,52 @@ async fn site_replication_targets_online(bucket: &str, replication_config_xml: &
true
}
/// Merge a peer's replication config into the local one.
///
/// `site-repl-*` rules encode the *sender's* outbound direction — their destination ARN
/// names the receiver — so applying a peer's rule set verbatim replaces the receiver's
/// reverse rule with one pointing at itself. No bucket target can satisfy that ARN
/// (`reconcile_site_replication_bucket_targets` skips the local peer), so the receiver
/// silently stops replicating back: the one-directional symptom. Only operator-authored
/// rules travel between sites; each site owns its own `site-repl-*` rules.
fn merge_incoming_replication_config(
incoming: Option<ReplicationConfiguration>,
local: Option<ReplicationConfiguration>,
) -> Option<ReplicationConfiguration> {
let incoming_role = incoming.as_ref().map(|config| config.role.clone()).unwrap_or_default();
// Operator rules first, then the local site rules — the same order
// `ensure_site_replication_bucket_replication_config_with_runtime` produces, so its
// no-op check matches and the bucket metadata is written once per broadcast, not twice.
let mut rules: Vec<ReplicationRule> = incoming
.into_iter()
.flat_map(|config| config.rules)
.filter(|rule| !is_site_replication_rule(rule))
.collect();
rules.extend(
local
.into_iter()
.flat_map(|config| config.rules)
.filter(is_site_replication_rule),
);
if rules.is_empty() {
return None;
}
for (index, rule) in rules.iter_mut().enumerate() {
rule.priority = Some(i32::try_from(index + 1).unwrap_or(i32::MAX));
}
// A site-replication ARN in `role` is the sender's, and `site_replication_target_arns_by_peer`
// reads it — carrying it over would pin the receiver's targets to the sender's identity.
let role = match replication_target_arn_deployment_id(&incoming_role) {
Some(_) => String::new(),
None => incoming_role,
};
Some(ReplicationConfiguration { role, rules })
}
/// Merge a peer's ILM expiry document into the local lifecycle config.
///
/// Mirrors MinIO's `mergeWithCurrentLCConfig` with one hardening: incoming
-1
View File
@@ -443,7 +443,6 @@ pub(crate) mod replication {
pub(crate) use super::ecstore_bucket::replication::{
REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION, REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS,
REPLICATION_CAPABILITY_CONTRACT_VERSION, REPLICATION_READ_ONLY_HISTORICAL_FIELDS, REPLICATION_WRITABLE_FIELDS,
is_site_replication_rule, merge_incoming_replication_config, replication_target_arn_deployment_id,
};
pub(crate) type BucketReplicationResyncStatus = super::ecstore_bucket::replication::BucketReplicationResyncStatus;
pub(crate) type BucketStats = super::ecstore_bucket::replication::BucketStats;
+13 -193
View File
@@ -38,9 +38,9 @@ use super::storage_api::bucket_usecase::bucket::{
metadata_sys,
policy_sys::PolicySys,
replication::{
ReplicationTargetValidationError, invalid_replication_config_status_field, is_site_replication_rule,
merge_incoming_replication_config, replication_target_arns, should_remove_replication_target,
unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns,
ReplicationTargetValidationError, invalid_replication_config_status_field, replication_target_arns,
should_remove_replication_target, unsupported_replication_config_field, validate_replication_config_structure,
validate_replication_config_target_arns,
},
target::{BucketTargetType, BucketTargets},
utils::serialize,
@@ -623,50 +623,11 @@ async fn validate_bucket_replication_update(bucket: &str, config: &ReplicationCo
validate_replication_config_targets(&targets, config)
}
/// Defense in depth for site-replication-managed buckets (issue #1948): an S3
/// PutBucketReplication replaces the operator-authored rules but must not wipe
/// the local `site-repl-*` rules the reconciler owns — until its next pass
/// (600s period) every peer link on this bucket would be silently dead. The
/// same merge also drops incoming `site-repl-*` impostor rules, matching the
/// peer bucket-meta ingestion path. Buckets without site-replication rules
/// keep the verbatim overwrite semantics.
fn merge_user_replication_config_update(
incoming: ReplicationConfiguration,
existing: Option<ReplicationConfiguration>,
) -> ReplicationConfiguration {
let has_site_rules = existing
.as_ref()
.is_some_and(|config| config.rules.iter().any(is_site_replication_rule));
if !has_site_rules {
return incoming;
}
// `existing` holds at least one site-replication rule the merge keeps, so
// the merged rule set is non-empty; the fallback only guards the type.
merge_incoming_replication_config(Some(incoming.clone()), existing).unwrap_or(incoming)
}
/// Split of an S3 DeleteBucketReplication on the stored config (issue #1948):
/// the operator-authored rules are removed, the local `site-repl-*` rules
/// survive (`None` means nothing survives and the config is deleted), and the
/// returned ARNs are the ones whose bucket targets may be garbage-collected —
/// never an ARN a surviving site-replication rule still points at.
fn split_replication_config_for_user_delete(
config: ReplicationConfiguration,
) -> (Option<ReplicationConfiguration>, HashSet<String>) {
let mut removable_arns = replication_target_arns(&config);
let remaining = merge_incoming_replication_config(None, Some(config));
if let Some(remaining) = remaining.as_ref() {
for rule in &remaining.rules {
removable_arns.remove(rule.destination.bucket.trim());
}
}
(remaining, removable_arns)
}
async fn replication_targets_without_arns(
async fn replication_targets_without_config_targets(
bucket: &str,
target_arns: &HashSet<String>,
config: &ReplicationConfiguration,
) -> S3Result<Option<(BucketTargets, usize)>> {
let target_arns = replication_target_arns(config);
if target_arns.is_empty() {
return Ok(None);
}
@@ -677,7 +638,7 @@ async fn replication_targets_without_arns(
Err(err) => return Err(ApiError::from(err).into()),
};
let removed = remove_replication_targets_from_config_targets(&mut targets, target_arns);
let removed = remove_replication_targets_from_config_targets(&mut targets, &target_arns);
if removed == 0 {
return Ok(None);
}
@@ -1643,29 +1604,15 @@ impl DefaultBucketUsecase {
Err(StorageError::ConfigNotFound) => None,
Err(err) => return Err(ApiError::from(err).into()),
};
let (remaining_config, updated_targets) = if let Some(config) = replication_config.as_ref() {
let (remaining, removable_arns) = split_replication_config_for_user_delete(config.clone());
let targets = replication_targets_without_arns(&bucket, &removable_arns).await?;
(remaining, targets)
let updated_targets = if let Some(config) = replication_config.as_ref() {
replication_targets_without_config_targets(&bucket, config).await?
} else {
(None, None)
None
};
match remaining_config {
// Site-replication rules and the targets backing them survive the
// S3 delete (issue #1948); only the operator-authored rules go.
Some(remaining) => {
let data = serialize_config(&remaining)?;
update_bucket_config_for_incarnation(&bucket, BUCKET_REPLICATION_CONFIG, data, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
}
None => {
delete_bucket_config_for_incarnation(&bucket, BUCKET_REPLICATION_CONFIG, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
}
}
delete_bucket_config_for_incarnation(&bucket, BUCKET_REPLICATION_CONFIG, expected_incarnation_id)
.await
.map_err(ApiError::from)?;
if let Some((targets, removed)) = updated_targets
&& let Err(err) =
write_replication_targets_after_config_delete(&bucket, &targets, removed, expected_incarnation_id).await
@@ -2538,12 +2485,6 @@ impl DefaultBucketUsecase {
let targets_guard = lock_bucket_targets_metadata(&bucket).await;
validate_bucket_replication_update(&bucket, &replication_configuration).await?;
let existing_config = match metadata_sys::get_replication_config(&bucket).await {
Ok((config, _)) => Some(config),
Err(StorageError::ConfigNotFound) => None,
Err(err) => return Err(ApiError::from(err).into()),
};
let replication_configuration = merge_user_replication_config_update(replication_configuration, existing_config);
let data = serialize_config(&replication_configuration)?;
update_bucket_config_for_incarnation(&bucket, BUCKET_REPLICATION_CONFIG, data, expected_incarnation_id)
.await
@@ -3173,127 +3114,6 @@ mod tests {
assert!(arns.contains(destination));
}
fn replication_rule_with_id(arn: &str, id: &str, priority: i32) -> ReplicationRule {
let mut rule = replication_rule_for_target(arn);
rule.id = Some(id.to_string());
rule.priority = Some(priority);
rule
}
#[test]
fn put_replication_merge_preserves_site_replication_rules() {
let existing = ReplicationConfiguration {
role: String::new(),
rules: vec![
replication_rule_with_id("arn:rustfs:replication::peer-dep:bucket", "site-repl-peer-dep", 1),
replication_rule_with_id("arn:rustfs:replication:us-east-1:old:bucket", "old-user-rule", 2),
],
};
let incoming = ReplicationConfiguration {
role: String::new(),
rules: vec![
replication_rule_with_id("arn:rustfs:replication:us-east-1:new:bucket", "new-user-rule", 1),
replication_rule_with_id("arn:rustfs:replication::forged-dep:bucket", "site-repl-forged", 2),
],
};
let merged = merge_user_replication_config_update(incoming, Some(existing));
let ids: Vec<_> = merged
.rules
.iter()
.map(|rule| rule.id.as_deref().unwrap_or_default())
.collect();
assert_eq!(
ids,
vec!["new-user-rule", "site-repl-peer-dep"],
"user rules replaced, local site-replication rule preserved, forged incoming site-repl rule dropped"
);
}
#[test]
fn put_replication_merge_returns_incoming_verbatim_without_site_rules() {
let existing = ReplicationConfiguration {
role: String::new(),
rules: vec![replication_rule_with_id(
"arn:rustfs:replication:us-east-1:old:bucket",
"old-user-rule",
7,
)],
};
let incoming = ReplicationConfiguration {
role: String::new(),
rules: vec![replication_rule_with_id(
"arn:rustfs:replication:us-east-1:new:bucket",
"new-user-rule",
5,
)],
};
let merged = merge_user_replication_config_update(incoming.clone(), Some(existing));
assert_eq!(merged.role, incoming.role);
assert_eq!(merged.rules, incoming.rules, "non-SR buckets keep the verbatim overwrite semantics");
}
#[test]
fn delete_replication_split_keeps_site_rules_and_their_targets() {
let sr_arn = "arn:rustfs:replication::peer-dep:bucket";
let user_arn = "arn:rustfs:replication:us-east-1:user:bucket";
let config = ReplicationConfiguration {
role: String::new(),
rules: vec![
replication_rule_with_id(user_arn, "user-rule", 1),
replication_rule_with_id(sr_arn, "site-repl-peer-dep", 2),
],
};
let (remaining, removable) = split_replication_config_for_user_delete(config);
let remaining = remaining.expect("site-replication rules must survive a user delete");
let ids: Vec<_> = remaining
.rules
.iter()
.map(|rule| rule.id.as_deref().unwrap_or_default())
.collect();
assert_eq!(ids, vec!["site-repl-peer-dep"]);
assert_eq!(removable, HashSet::from([user_arn.to_string()]));
}
#[test]
fn delete_replication_split_protects_targets_shared_with_site_rules() {
let sr_arn = "arn:rustfs:replication::peer-dep:bucket";
let config = ReplicationConfiguration {
role: String::new(),
rules: vec![
replication_rule_with_id(sr_arn, "user-rule-on-sr-target", 1),
replication_rule_with_id(sr_arn, "site-repl-peer-dep", 2),
],
};
let (remaining, removable) = split_replication_config_for_user_delete(config);
assert!(remaining.is_some());
assert!(
removable.is_empty(),
"a target still referenced by a surviving site-replication rule must not be removed"
);
}
#[test]
fn delete_replication_split_removes_everything_without_site_rules() {
let user_arn = "arn:rustfs:replication:us-east-1:user:bucket";
let config = ReplicationConfiguration {
role: String::new(),
rules: vec![replication_rule_with_id(user_arn, "user-rule", 1)],
};
let (remaining, removable) = split_replication_config_for_user_delete(config);
assert!(remaining.is_none(), "without site-replication rules the whole config is deleted");
assert_eq!(removable, HashSet::from([user_arn.to_string()]));
}
fn replication_targets_with_arn(arns: &[&str]) -> BucketTargets {
BucketTargets {
targets: arns
-2
View File
@@ -614,8 +614,6 @@ pub(crate) mod bucket {
use crate::storage::storage_api::ecstore_bucket::replication as replication_contracts;
pub(crate) use replication_contracts::{is_site_replication_rule, merge_incoming_replication_config};
type ReplicationObjectBridge = crate::storage::storage_api::ecstore_bucket::replication::ReplicationObjectBridge;
pub(crate) type DeleteReplicationConfigSnapshot =
crate::storage::storage_api::ecstore_bucket::replication::DeleteReplicationConfigSnapshot;
-150
View File
@@ -63,54 +63,6 @@ use crate::app::storage_api::object_usecase::bucket::replication::{
};
use crate::storage::storage_api::ecfs_consumer::StorageObjectOptions as ObjectOptions;
#[cfg(test)]
static SITE_REPLICATION_GATE_TEST_OVERRIDE: std::sync::atomic::AtomicU8 = std::sync::atomic::AtomicU8::new(0);
#[cfg(test)]
const SITE_REPLICATION_GATE_FORCE_DISABLED: u8 = 1;
#[cfg(test)]
const SITE_REPLICATION_GATE_FORCE_ENABLED: u8 = 2;
async fn site_replication_gate_enabled() -> S3Result<bool> {
#[cfg(test)]
match SITE_REPLICATION_GATE_TEST_OVERRIDE.load(std::sync::atomic::Ordering::SeqCst) {
SITE_REPLICATION_GATE_FORCE_DISABLED => return Ok(false),
SITE_REPLICATION_GATE_FORCE_ENABLED => return Ok(true),
_ => {}
}
crate::admin::handlers::site_replication::site_replication_enabled().await
}
/// MinIO `ErrReplicationDenyEditError`.
fn replication_deny_edit_error() -> S3Error {
let mut err = S3Error::with_message(
S3ErrorCode::Custom("XMinioReplicationDenyEdit".into()),
"Sub-User is not allowed to edit Replication configuration",
);
err.set_status_code(StatusCode::BAD_REQUEST);
err
}
/// Site-replication gate for S3 replication-config edits (issue #1948).
///
/// On a site-replication deployment the bucket's replication config carries
/// the operator-managed `site-repl-*` rules that keep every peer in sync, and
/// a successful edit is broadcast to all peers — so a user holding only
/// bucket-scoped `s3:PutReplicationConfiguration` could rewrite or erase
/// replication net-wide. MinIO parity (`ErrReplicationDenyEditError`): only
/// owner credentials (root or root-parented) may edit. Runs after the policy
/// authorization in the access layer and only on the external S3 path — the
/// reconciler and peer bucket-meta ingestion never route through these
/// handlers.
async fn deny_replication_config_edit_for_non_owner<T>(req: &S3Request<T>) -> S3Result<()> {
if crate::storage::access::req_info_ref(req)?.is_owner {
return Ok(());
}
if site_replication_gate_enabled().await? {
return Err(replication_deny_edit_error());
}
Ok(())
}
#[derive(Debug, Clone)]
pub struct FS {
/// This server's late-bound application-context slot (backlog#1052 S2).
@@ -548,7 +500,6 @@ impl S3 for FS {
&self,
req: S3Request<DeleteBucketReplicationInput>,
) -> S3Result<S3Response<DeleteBucketReplicationOutput>> {
deny_replication_config_edit_for_non_owner(&req).await?;
let usecase = s3_api::bucket_usecase_for(self);
usecase.execute_delete_bucket_replication(req).await
}
@@ -1402,7 +1353,6 @@ impl S3 for FS {
&self,
req: S3Request<PutBucketReplicationInput>,
) -> S3Result<S3Response<PutBucketReplicationOutput>> {
deny_replication_config_edit_for_non_owner(&req).await?;
let usecase = s3_api::bucket_usecase_for(self);
usecase.execute_put_bucket_replication(req).await
}
@@ -1969,103 +1919,3 @@ impl S3 for FS {
Box::pin(usecase.execute_upload_part_copy(req)).await
}
}
#[cfg(test)]
mod tests {
use super::{
FS, SITE_REPLICATION_GATE_FORCE_DISABLED, SITE_REPLICATION_GATE_FORCE_ENABLED, SITE_REPLICATION_GATE_TEST_OVERRIDE,
};
use crate::storage::access::ReqInfo;
use http::Method;
use http::StatusCode;
use s3s::dto::{DeleteBucketReplicationInput, PutBucketReplicationInput, ReplicationConfiguration};
use s3s::{S3, S3Error, S3ErrorCode, S3Request};
use std::sync::atomic::Ordering;
fn replication_config_edit_request<T>(input: T, is_owner: bool) -> S3Request<T> {
let mut req = S3Request {
input,
method: Method::PUT,
uri: http::Uri::from_static("/"),
headers: http::HeaderMap::new(),
extensions: http::Extensions::new(),
credentials: None,
region: None,
service: None,
trailing_headers: None,
};
req.extensions.insert(ReqInfo {
is_owner,
..Default::default()
});
req
}
fn put_bucket_replication_input() -> PutBucketReplicationInput {
PutBucketReplicationInput {
bucket: "test-bucket".to_string(),
checksum_algorithm: None,
content_md5: None,
expected_bucket_owner: None,
replication_configuration: ReplicationConfiguration {
role: String::new(),
rules: Vec::new(),
},
token: None,
}
}
fn delete_bucket_replication_input() -> DeleteBucketReplicationInput {
DeleteBucketReplicationInput {
bucket: "test-bucket".to_string(),
expected_bucket_owner: None,
}
}
fn assert_replication_deny_edit(err: &S3Error) {
match err.code() {
S3ErrorCode::Custom(code) => assert_eq!(code, "XMinioReplicationDenyEdit"),
other => panic!("expected XMinioReplicationDenyEdit, got {other:?}"),
}
assert_eq!(err.status_code(), Some(StatusCode::BAD_REQUEST));
}
/// Single test on purpose: the branches share the process-wide gate
/// override, and parallel tests would race it.
#[tokio::test]
async fn replication_config_edit_gate_denies_only_non_owner_under_site_replication() {
let fs = FS::new();
SITE_REPLICATION_GATE_TEST_OVERRIDE.store(SITE_REPLICATION_GATE_FORCE_ENABLED, Ordering::SeqCst);
// Non-owner PUT/DELETE through the real S3 handlers: denied by the
// gate before the usecase (and thus the store) is ever touched.
let err = fs
.put_bucket_replication(replication_config_edit_request(put_bucket_replication_input(), false))
.await
.expect_err("non-owner PutBucketReplication must be denied while site replication is enabled");
assert_replication_deny_edit(&err);
let err = fs
.delete_bucket_replication(replication_config_edit_request(delete_bucket_replication_input(), false))
.await
.expect_err("non-owner DeleteBucketReplication must be denied while site replication is enabled");
assert_replication_deny_edit(&err);
// Owner passes the gate (the usecase's empty-rules structure error
// proves the request reached the usecase instead of the deny path).
let err = fs
.put_bucket_replication(replication_config_edit_request(put_bucket_replication_input(), true))
.await
.expect_err("owner request should pass the gate and fail later on config validation");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
// Without site replication the policy check alone still governs the edit.
SITE_REPLICATION_GATE_TEST_OVERRIDE.store(SITE_REPLICATION_GATE_FORCE_DISABLED, Ordering::SeqCst);
let err = fs
.put_bucket_replication(replication_config_edit_request(put_bucket_replication_input(), false))
.await
.expect_err("non-owner request should pass the gate and fail later on config validation");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
SITE_REPLICATION_GATE_TEST_OVERRIDE.store(0, Ordering::SeqCst);
}
}
+5 -4
View File
@@ -554,10 +554,11 @@ fn apply_replication_timestamps_from_headers(headers: &HeaderMap<HeaderValue>, o
// Persist into the internal metadata keys so a later outbound replication
// pass (replication_target_boundary) reads the source's modification
// times instead of falling back to mod_time.
// TODO(P1-6): receiver-side LWW is still missing — when the stored
// per-category timestamp is newer than the inbound one, the existing
// tags/retention/legal-hold should win instead of being overwritten.
// times instead of falling back to mod_time. Receiver-side LWW happens at
// the set layer under the object write lock
// (ecstore set_disk::ops::object::merge_replication_metadata_lww,
// rustfs/backlog#1953): a category whose stored timestamp is newer than
// the inbound one keeps the local values.
for (timestamp, suffix) in [
(opts.replication_tagging_timestamp, SUFFIX_TAGGING_TIMESTAMP),
(opts.replication_retention_timestamp, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP),