From 31933c32f9e925f73c4654e3268f5cb594cdf2b8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Sun, 23 Aug 2026 20:17:50 +0800 Subject: [PATCH] fix(replication): apply receiver-side LWW to inbound metadata categories (#6379) --- crates/e2e_test/src/lib.rs | 5 + .../src/replication_extension_test.rs | 422 ++++++++++++++- .../src/replication_lww_receiver_test.rs | 155 ++++++ .../ecstore/src/bucket/bucket_target_sys.rs | 43 +- .../bucket/replication/replication_pool.rs | 24 + .../replication/replication_resyncer.rs | 140 ++++- .../replication_target_boundary.rs | 32 +- crates/ecstore/src/set_disk/ops/multipart.rs | 177 ++++++ crates/ecstore/src/set_disk/ops/object.rs | 509 ++++++++++++++++++ crates/replication/src/operation.rs | 18 + crates/utils/src/http/headers.rs | 2 + rustfs/src/storage/options.rs | 9 +- 12 files changed, 1523 insertions(+), 13 deletions(-) create mode 100644 crates/e2e_test/src/replication_lww_receiver_test.rs diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index b3c1a175a..28c63b267 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -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; diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 3d51d139b..4d7b641e4 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -57,7 +57,7 @@ use rustfs_madmin::{ AddServiceAccountReq, ListServiceAccountsResp, PeerInfo, PeerSite, ReplicateAddStatus, ReplicateEditStatus, ReplicateRemoveStatus, SRRemoveReq, SRResyncOpStatus, SRStatusInfo, SiteReplicationInfo, SyncStatus, }; -use s3s::header::X_AMZ_REPLICATION_STATUS; +use s3s::header::{X_AMZ_REPLICATION_STATUS, X_AMZ_TAGGING}; use sha2::{Digest, Sha256}; use std::collections::BTreeMap; use std::convert::Infallible; @@ -2023,6 +2023,7 @@ async fn forward_replication_proxy_request( client: &reqwest::Client, request_count: &AtomicU64, mut replication_enabled: watch::Receiver, + mut held_tagging: watch::Receiver>, ) -> Response> { let (parts, body) = request.into_parts(); let is_replication = parts @@ -2036,6 +2037,17 @@ async fn forward_replication_proxy_request( return proxy_error_response("replication gate closed"); } } + // Content-keyed hold: park only the replication request whose + // `x-amz-tagging` matches the held value, letting every other delivery + // through, so a test can make one specific (stale) delivery the last + // write the backend sees. + if let Some(tagging) = parts.headers.get(X_AMZ_TAGGING).and_then(|value| value.to_str().ok()) { + while held_tagging.borrow().as_deref() == Some(tagging) { + if held_tagging.changed().await.is_err() { + return proxy_error_response("replication tag hold closed"); + } + } + } } let Some(path_and_query) = parts.uri.path_and_query() else { @@ -2070,12 +2082,26 @@ async fn start_replication_counting_proxy( backend_url: &str, tasks: &mut JoinSet<()>, ) -> Result<(String, Arc, watch::Sender), Box> { + let (proxy_url, request_count, replication_enabled, _held_tagging) = + start_replication_counting_proxy_with_tag_hold(backend_url, tasks).await?; + Ok((proxy_url, request_count, replication_enabled)) +} + +/// [`start_replication_counting_proxy`] plus a content-keyed hold: while the +/// returned `watch::Sender>` holds `Some(tagging)`, replication +/// requests whose `x-amz-tagging` equals `tagging` are parked (and still +/// counted); all other traffic flows. Send `None` to release them. +async fn start_replication_counting_proxy_with_tag_hold( + backend_url: &str, + tasks: &mut JoinSet<()>, +) -> Result<(String, Arc, watch::Sender, watch::Sender>), Box> { let listener = TcpListener::bind("127.0.0.1:0").await?; let proxy_url = format!("http://{}", listener.local_addr()?); let backend_url = backend_url.to_string(); let request_count = Arc::new(AtomicU64::new(0)); let task_request_count = request_count.clone(); let (replication_enabled, task_replication_enabled) = watch::channel(true); + let (held_tagging, task_held_tagging) = watch::channel(None); tasks.spawn(async move { let client = local_http_client(); let mut connections = JoinSet::new(); @@ -2087,12 +2113,14 @@ async fn start_replication_counting_proxy( let client = client.clone(); let request_count = task_request_count.clone(); let replication_enabled = task_replication_enabled.clone(); + let held_tagging = task_held_tagging.clone(); connections.spawn(async move { let service = service_fn(move |request| { let backend_url = backend_url.clone(); let client = client.clone(); let request_count = request_count.clone(); let replication_enabled = replication_enabled.clone(); + let held_tagging = held_tagging.clone(); async move { Ok::<_, Infallible>( forward_replication_proxy_request( @@ -2101,6 +2129,7 @@ async fn start_replication_counting_proxy( &client, &request_count, replication_enabled, + held_tagging, ) .await, ) @@ -2113,7 +2142,7 @@ async fn start_replication_counting_proxy( } } }); - Ok((proxy_url, request_count, replication_enabled)) + Ok((proxy_url, request_count, replication_enabled, held_tagging)) } async fn site_replication_remove( @@ -6949,6 +6978,395 @@ async fn test_site_replication_active_active_converges_without_loops_real_dual_n } } +/// Replication status a site reports for one object version via HEAD +/// (`x-amz-replication-status`), or `None` when the header is absent. +async fn head_replication_status( + client: &Client, + bucket: &str, + key: &str, + version_id: &str, +) -> Result, Box> { + let head = client + .head_object() + .bucket(bucket) + .key(key) + .version_id(version_id) + .send() + .await?; + Ok(head.replication_status().map(|status| status.as_str().to_string())) +} + +/// Poll one site until the version's replication status is one of `expected`. +async fn wait_for_version_replication_status( + client: &Client, + bucket: &str, + key: &str, + version_id: &str, + expected: &[&str], + site: &str, +) -> Result> { + let deadline = tokio::time::Instant::now() + Duration::from_secs(60); + loop { + let last = head_replication_status(client, bucket, key, version_id).await?; + if let Some(status) = last.as_deref() + && expected.contains(&status) + { + return Ok(status.to_string()); + } + if tokio::time::Instant::now() >= deadline { + return Err(format!( + "{site}: {bucket}/{key}?versionId={version_id} replication status {last:?} never reached {expected:?}" + ) + .into()); + } + sleep(Duration::from_millis(200)).await; + } +} + +async fn put_single_tag( + client: &Client, + bucket: &str, + key: &str, + version_id: &str, + tag_key: &str, + tag_value: &str, +) -> Result<(), Box> { + client + .put_object_tagging() + .bucket(bucket) + .key(key) + .version_id(version_id) + .tagging( + aws_sdk_s3::types::Tagging::builder() + .tag_set(aws_sdk_s3::types::Tag::builder().key(tag_key).value(tag_value).build()?) + .build()?, + ) + .send() + .await?; + Ok(()) +} + +async fn get_single_tag( + client: &Client, + bucket: &str, + key: &str, + version_id: &str, + tag_key: &str, +) -> Result, Box> { + let tagging = client + .get_object_tagging() + .bucket(bucket) + .key(key) + .version_id(version_id) + .send() + .await?; + Ok(tagging + .tag_set() + .iter() + .find(|tag| tag.key() == tag_key) + .map(|tag| tag.value().to_string())) +} + +/// Poll one site until the version's `tag_key` equals `expected`. +async fn wait_for_single_tag( + client: &Client, + bucket: &str, + key: &str, + version_id: &str, + tag_key: &str, + expected: &str, + site: &str, +) -> Result<(), Box> { + let deadline = tokio::time::Instant::now() + Duration::from_secs(60); + loop { + let observed = get_single_tag(client, bucket, key, version_id, tag_key).await?; + if observed.as_deref() == Some(expected) { + return Ok(()); + } + if tokio::time::Instant::now() >= deadline { + return Err(format!( + "{site}: {bucket}/{key}?versionId={version_id} tag {tag_key}={observed:?} never became {expected}" + ) + .into()); + } + sleep(Duration::from_millis(200)).await; + } +} + +/// Tag key the dual-node LWW scenario edits on both sites. +const LWW_TAG_KEY: &str = "owner"; + +/// Assert the version's [`LWW_TAG_KEY`] stays `expected` on both sites for a +/// full quiet window (no late stale delivery flips it back). +async fn assert_tag_stable_on_both_sites( + site_a_client: &Client, + site_b_client: &Client, + bucket: &str, + key: &str, + version_id: &str, + expected: &str, + quiet: Duration, +) -> Result<(), Box> { + let deadline = tokio::time::Instant::now() + quiet; + loop { + let on_a = get_single_tag(site_a_client, bucket, key, version_id, LWW_TAG_KEY).await?; + let on_b = get_single_tag(site_b_client, bucket, key, version_id, LWW_TAG_KEY).await?; + assert_eq!(on_a.as_deref(), Some(expected), "site A tag {LWW_TAG_KEY} regressed from the LWW winner"); + assert_eq!(on_b.as_deref(), Some(expected), "site B tag {LWW_TAG_KEY} regressed from the LWW winner"); + if tokio::time::Instant::now() >= deadline { + return Ok(()); + } + sleep(Duration::from_millis(250)).await; + } +} + +/// Wait until the counting proxy in front of a site has admitted `expected` +/// replication requests in total (requests held by a closed gate still count). +async fn wait_for_proxy_replication_requests( + counter: &AtomicU64, + expected: u64, + site: &str, +) -> Result<(), Box> { + let deadline = tokio::time::Instant::now() + Duration::from_secs(60); + loop { + let observed = counter.load(Ordering::Relaxed); + if observed >= expected { + return Ok(()); + } + if tokio::time::Instant::now() >= deadline { + return Err(format!("{site} proxy saw {observed} replication requests, expected at least {expected}").into()); + } + sleep(Duration::from_millis(25)).await; + } +} + +/// rustfs/backlog#1953 (audit A4/P1-6): receiver-side LWW for replicated +/// metadata categories, exercised end to end over the real dual-node +/// active-active site-replication control plane — sender, worker, status +/// bookkeeping and persisted failure recovery all participate (the single-server +/// `replication_lww_receiver_test` only injects authorized replication PUTs). +/// +/// Scenario on one versioned object: +/// 1. reciprocal tag edits in real order (A then B) converge both sites on the +/// newer tag and leave the author COMPLETED / the receiver REPLICA; +/// 2. out-of-order delivery: A's edit is held at B's inbound proxy while B +/// authors a newer edit that reaches A first; releasing the stale delivery +/// must NOT roll B back — both sites settle on B's value and stay there +/// through a quiet window, with no FAILED/PENDING status left behind; +/// 3. persisted retry: B is stopped, A's delivery reaches FAILED, A restarts, +/// then B returns and the scanner-replayed edit converges both sites forward. +/// Durable metadata-MRF serialization/reconstruction is covered separately by +/// `metadata_mrf_roundtrip_preserves_tags_and_admitted_targets`. +#[tokio::test] +async fn test_site_replication_tagging_lww_converges_active_active_real_dual_node() -> TestResult { + init_logging(); + + match tokio::time::timeout(Duration::from_secs(420), async { + // The scanner is fast for the final persisted-failure recovery phase. + // Step 2 finishes and proves a quiet stable winner before that phase, + // so a later scanner pass cannot mask its stale-delivery assertion. + let mut site_env = replication_fast_env(); + site_env.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV); + site_env.extend_from_slice(FAST_SCANNER_ENV); + + let mut site_a_env = RustFSTestEnvironment::new().await?; + site_a_env.start_rustfs_server_with_env(vec![], &site_env).await?; + + let mut site_b_env = RustFSTestEnvironment::new().await?; + site_b_env.start_rustfs_server_with_env(vec![], &site_env).await?; + + let mut proxy_tasks = JoinSet::new(); + let (site_a_proxy, site_a_replication_requests, _site_a_replication_enabled, site_a_held_tagging) = + start_replication_counting_proxy_with_tag_hold(&site_a_env.url, &mut proxy_tasks).await?; + let (site_b_proxy, site_b_replication_requests, _site_b_replication_enabled, site_b_held_tagging) = + start_replication_counting_proxy_with_tag_hold(&site_b_env.url, &mut proxy_tasks).await?; + + let site_a_client = site_a_env.create_s3_client(); + let site_b_client = site_b_env.create_s3_client(); + let bucket = "site-repl-tag-lww"; + let key = "lww.txt"; + + let add_status = site_replication_add( + &site_a_env, + &[ + PeerSite { + name: "lww-site-a".to_string(), + endpoint: site_a_env.url.clone(), + access_key: site_a_env.access_key.clone(), + secret_key: site_a_env.secret_key.clone(), + ..Default::default() + }, + PeerSite { + name: "lww-site-b".to_string(), + endpoint: site_b_env.url.clone(), + access_key: site_b_env.access_key.clone(), + secret_key: site_b_env.secret_key.clone(), + ..Default::default() + }, + ], + ) + .await?; + assert!(add_status.success, "unexpected site add result: {add_status:?}"); + + let site_info = wait_for_site_replication_enabled(&site_a_env, 2).await?; + wait_for_site_replication_enabled(&site_b_env, 2).await?; + + // Route both directions through the counting proxies so inbound + // replication to B can be held (out-of-order delivery) and observed. + for (env_url, proxy_url, label) in [(&site_a_env.url, &site_a_proxy, "A"), (&site_b_env.url, &site_b_proxy, "B")] { + let mut peer = site_info + .sites + .iter() + .find(|peer| peer.endpoint == *env_url) + .ok_or_else(|| format!("site {label} peer missing from replication info"))? + .clone(); + peer.endpoint = proxy_url.clone(); + peer.sync_state = SyncStatus::Enable; + let edit = site_replication_edit(&site_a_env, "", &peer).await?; + assert!(edit.success, "unexpected site {label} endpoint edit: {edit:?}"); + } + for env in [&site_a_env, &site_b_env] { + wait_for_site_replication_info(env, |info| { + info.sites.iter().any(|peer| peer.endpoint == site_a_proxy) + && info.sites.iter().any(|peer| peer.endpoint == site_b_proxy) + }) + .await?; + } + + site_a_client.create_bucket().bucket(bucket).send().await?; + wait_for_bucket_on_target(&site_b_client, bucket).await?; + + let version_id = site_a_client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"tag lww payload")) + .send() + .await? + .version_id() + .ok_or("site A PUT omitted version ID")? + .to_string(); + wait_for_replicated_object(&site_b_client, bucket, key, "tag lww payload").await?; + wait_for_version_replication_status(&site_a_client, bucket, key, &version_id, &["COMPLETED"], "site A").await?; + wait_for_version_replication_status(&site_b_client, bucket, key, &version_id, &["REPLICA"], "site B").await?; + wait_for_proxy_replication_requests(&site_b_replication_requests, 1, "site B").await?; + + // --- 1. reciprocal edits in real order: A then B ---------------------- + put_single_tag(&site_a_client, bucket, key, &version_id, LWW_TAG_KEY, "a1").await?; + wait_for_single_tag(&site_b_client, bucket, key, &version_id, LWW_TAG_KEY, "a1", "site B").await?; + wait_for_version_replication_status(&site_a_client, bucket, key, &version_id, &["COMPLETED"], "site A").await?; + wait_for_version_replication_status(&site_b_client, bucket, key, &version_id, &["REPLICA"], "site B").await?; + wait_for_proxy_replication_requests(&site_b_replication_requests, 2, "site B").await?; + + put_single_tag(&site_b_client, bucket, key, &version_id, LWW_TAG_KEY, "b1").await?; + wait_for_single_tag(&site_a_client, bucket, key, &version_id, LWW_TAG_KEY, "b1", "site A").await?; + wait_for_version_replication_status(&site_b_client, bucket, key, &version_id, &["COMPLETED"], "site B").await?; + wait_for_version_replication_status(&site_a_client, bucket, key, &version_id, &["REPLICA"], "site A").await?; + assert_tag_stable_on_both_sites(&site_a_client, &site_b_client, bucket, key, &version_id, "b1", Duration::from_secs(3)) + .await?; + + // --- 2. concurrent edits, stale delivery last ------------------------ + // Both sites edit the same version while each other's delivery is + // parked at the peer's inbound proxy (content-keyed: only the + // `owner=a2` / `owner=b2` replication PUTs wait, everything else + // flows). B's edit is the newer one. Releasing A's stale `a2` first + // makes it the last write B sees while A itself still holds `a2`, so + // nothing A could re-deliver carries the winner: only receiver-side + // LWW on B can keep `b2`. Releasing `b2` afterwards converges A. + site_b_held_tagging.send(Some("owner=a2".to_string()))?; + site_a_held_tagging.send(Some("owner=b2".to_string()))?; + let a2_parked_at = site_b_replication_requests.load(Ordering::Relaxed) + 1; + let b2_parked_at = site_a_replication_requests.load(Ordering::Relaxed) + 1; + put_single_tag(&site_a_client, bucket, key, &version_id, LWW_TAG_KEY, "a2").await?; + wait_for_proxy_replication_requests(&site_b_replication_requests, a2_parked_at, "site B").await?; + sleep(Duration::from_millis(50)).await; + put_single_tag(&site_b_client, bucket, key, &version_id, LWW_TAG_KEY, "b2").await?; + wait_for_proxy_replication_requests(&site_a_replication_requests, b2_parked_at, "site A").await?; + assert_eq!( + get_single_tag(&site_a_client, bucket, key, &version_id, LWW_TAG_KEY) + .await? + .as_deref(), + Some("a2") + ); + assert_eq!( + get_single_tag(&site_b_client, bucket, key, &version_id, LWW_TAG_KEY) + .await? + .as_deref(), + Some("b2") + ); + + // Release the stale a2 delivery onto B: the newer local b2 must + // survive, and the delivery itself must still succeed (A reaches + // COMPLETED instead of looping through MRF with the stale value). + site_b_held_tagging.send(None)?; + wait_for_version_replication_status(&site_a_client, bucket, key, &version_id, &["COMPLETED"], "site A").await?; + let stale_deadline = tokio::time::Instant::now() + Duration::from_secs(3); + loop { + assert_eq!( + get_single_tag(&site_b_client, bucket, key, &version_id, LWW_TAG_KEY) + .await? + .as_deref(), + Some("b2"), + "a stale inbound delivery rolled back site B's newer tag (receiver-side LWW regression)" + ); + if tokio::time::Instant::now() >= stale_deadline { + break; + } + sleep(Duration::from_millis(250)).await; + } + + // Release b2 onto A: the newer edit wins there and both sites settle. + // B's own version may legitimately read REPLICA here: the stale inbound + // a2 write re-labelled it as a replica write (keeping B's tags); what + // must not remain is PENDING/FAILED. + site_a_held_tagging.send(None)?; + wait_for_single_tag(&site_a_client, bucket, key, &version_id, LWW_TAG_KEY, "b2", "site A").await?; + wait_for_version_replication_status(&site_b_client, bucket, key, &version_id, &["COMPLETED", "REPLICA"], "site B") + .await?; + assert_tag_stable_on_both_sites(&site_a_client, &site_b_client, bucket, key, &version_id, "b2", Duration::from_secs(4)) + .await?; + for (client, site) in [(&site_a_client, "site A"), (&site_b_client, "site B")] { + let status = head_replication_status(client, bucket, key, &version_id).await?; + assert!( + matches!(status.as_deref(), Some("COMPLETED" | "REPLICA")), + "{site} must not be left PENDING/FAILED after the concurrent edits: {status:?}" + ); + } + + // --- 3. persisted FAILED state survives a source restart ------------ + site_b_env.stop_server(); + put_single_tag(&site_a_client, bucket, key, &version_id, LWW_TAG_KEY, "a3").await?; + wait_for_version_replication_status(&site_a_client, bucket, key, &version_id, &["FAILED"], "site A").await?; + site_a_env.restart_server_preserving_data(vec![], &site_env).await?; + wait_for_site_replication_enabled(&site_a_env, 2).await?; + site_b_env.restart_server_preserving_data(vec![], &site_env).await?; + wait_for_site_replication_enabled(&site_b_env, 2).await?; + + wait_for_single_tag(&site_b_client, bucket, key, &version_id, LWW_TAG_KEY, "a3", "site B").await?; + wait_for_version_replication_status(&site_a_client, bucket, key, &version_id, &["COMPLETED"], "site A").await?; + wait_for_version_replication_status(&site_b_client, bucket, key, &version_id, &["REPLICA"], "site B").await?; + assert_tag_stable_on_both_sites(&site_a_client, &site_b_client, bucket, key, &version_id, "a3", Duration::from_secs(3)) + .await?; + + // The object itself never forked: one version on each side. + tokio::time::timeout( + Duration::from_secs(70), + assert_replication_converged(&site_a_client, bucket, &site_b_client, bucket), + ) + .await??; + let state = list_replication_state(&site_a_client, bucket).await?; + assert_eq!(state.len(), 1, "tag edits must not create new object versions: {state:?}"); + assert_eq!(state[0].version_id, version_id); + + proxy_tasks.abort_all(); + Ok(()) + }) + .await + { + Ok(result) => result, + Err(_) => Err("site replication tagging LWW test timed out".into()), + } +} #[tokio::test] async fn test_site_replication_replicates_policy_backed_user_access_real_dual_node() -> Result<(), Box> { init_logging(); diff --git a/crates/e2e_test/src/replication_lww_receiver_test.rs b/crates/e2e_test/src/replication_lww_receiver_test.rs new file mode 100644 index 000000000..4c856213c --- /dev/null +++ b/crates/e2e_test/src/replication_lww_receiver_test.rs @@ -0,0 +1,155 @@ +#![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). +//! +//! The real dual-site path (sender, worker, status bookkeeping, MRF replay) +//! is covered by +//! `replication_extension_test::test_site_replication_tagging_lww_converges_active_active_real_dual_node`; +//! this file stays as the fast, single-process receiver check. + +use crate::common::{RustFSTestEnvironment, init_logging}; +use aws_sdk_s3::Client; +use aws_sdk_s3::primitives::ByteStream; +use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration}; + +type TestResult = Result<(), Box>; + +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 { + 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")] +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(()) +} diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index 7e5a5df5f..4c4e66a79 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -58,8 +58,8 @@ use rustfs_config::{DEFAULT_TRUST_LEAF_CERT_AS_CA, ENV_TRUST_LEAF_CERT_AS_CA, RU use rustfs_utils::egress::{OutboundUrlError, validate_outbound_url}; use rustfs_utils::http::{ AMZ_BUCKET_REPLICATION_STATUS, AMZ_OBJECT_LOCK_BYPASS_GOVERNANCE, AMZ_OBJECT_LOCK_LEGAL_HOLD, AMZ_OBJECT_LOCK_MODE, - AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, AMZ_STORAGE_CLASS, AMZ_WEBSITE_REDIRECT_LOCATION, is_amz_header, is_minio_header, - is_rustfs_header, is_standard_header, is_storageclass_header, + AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, AMZ_OBJECT_TAGGING_LOWER, AMZ_STORAGE_CLASS, AMZ_WEBSITE_REDIRECT_LOCATION, is_amz_header, + is_minio_header, is_rustfs_header, is_standard_header, is_storageclass_header, }; use rustfs_utils::http::{ SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_PROXY_REQUEST, @@ -1774,6 +1774,22 @@ impl PutObjectOptions { Self::insert_checked(&mut header, AMZ_BUCKET_REPLICATION_STATUS, self.internal.replication_status.as_str()); } + // MinIO PutObjectOptions.Header parity: object tags travel on the + // `x-amz-tagging` header (form-urlencoded). `replication_put_object_options` + // fills `user_tags` from the source version; without this header the + // whole-object transport delivered a tagless replica, so tag edits + // never reached the peer and the receiver-side LWW comparison + // (rustfs/backlog#1953) had nothing to judge. + if !self.user_tags.is_empty() { + let mut tags: Vec<(&String, &String)> = self.user_tags.iter().collect(); + tags.sort(); + let mut encoded = url::form_urlencoded::Serializer::new(String::new()); + for (key, value) in tags { + encoded.append_pair(key, value); + } + Self::insert_checked(&mut header, AMZ_OBJECT_TAGGING_LOWER, &encoded.finish()); + } + for (k, v) in &self.user_metadata { let Ok(header_value) = HeaderValue::from_str(v) else { warn!("skipping user metadata header with invalid value: {}", k); @@ -3195,6 +3211,29 @@ mod tests { } } + #[test] + fn put_object_headers_carry_user_tags_on_x_amz_tagging() { + // rustfs/backlog#1953: tag edits replicate through the whole-object + // transport, so the source tags must travel on x-amz-tagging. + let mut opts = PutObjectOptions::default(); + opts.user_tags.insert("owner".to_string(), "site a".to_string()); + opts.user_tags.insert("env".to_string(), "prod".to_string()); + + let header = opts.header(); + let tagging = header + .get(AMZ_OBJECT_TAGGING_LOWER) + .expect("user tags must be transported on x-amz-tagging") + .to_str() + .expect("tag header must be ASCII"); + // Deterministic key order; values are form-urlencoded. + assert_eq!(tagging, "env=prod&owner=site+a"); + + assert!( + PutObjectOptions::default().header().get(AMZ_OBJECT_TAGGING_LOWER).is_none(), + "a tagless source must not send an empty x-amz-tagging header" + ); + } + #[test] fn put_object_headers_omit_unset_replication_timestamps() { // UNIX_EPOCH means "never modified on the source"; sending it would diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 61dbf2f1e..42934318a 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -4194,6 +4194,30 @@ mod tests { assert_eq!(ri.checksum, Some(checksum)); } + #[test] + fn metadata_mrf_roundtrip_preserves_tags_and_admitted_targets() { + let target = "arn:rustfs:replication:target-a"; + let object = ObjectInfo { + bucket: "source".to_string(), + name: "object".to_string(), + version_id: Some(Uuid::new_v4()), + user_tags: Arc::new("owner=a3".to_string()), + ..Default::default() + }; + let live = + replicate_object_info_from_object_info(object.clone(), test_replicate_decision(&[target]), ReplicationType::Metadata); + let persisted = live.to_mrf_entry(); + let encoded = encode_mrf_file(std::slice::from_ref(&persisted)).expect("metadata MRF entry should encode"); + let decoded = decode_mrf_file(&encoded).expect("metadata MRF entry should decode"); + + assert_eq!(decoded[0].op, MrfOpKind::Metadata); + assert_eq!(decoded[0].target_arns, vec![target.to_string()]); + let replayed = admitted_mrf_replicate_object(object, &decoded[0], ReplicationType::Metadata); + assert_eq!(replayed.op_type, ReplicationType::Metadata); + assert_eq!(replayed.user_tags, "owner=a3"); + assert_eq!(replayed.admitted_target_arns(), vec![target.to_string()]); + } + #[tokio::test] async fn mrf_save_admission_waits_for_capacity_instead_of_dropping() { let (tx, mut rx) = mpsc::channel(1); diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index e670a7d6f..587be7ae3 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -76,7 +76,8 @@ use metrics::counter; use rmp_serde; use rustfs_s3_types::EventName; use rustfs_utils::http::{ - AMZ_TAGGING_DIRECTIVE, SUFFIX_REPLICATION_RESET, SUFFIX_REPLICATION_STATUS, has_internal_suffix, insert_str, + AMZ_BUCKET_REPLICATION_STATUS, AMZ_TAGGING_DIRECTIVE, SUFFIX_REPLICATION_RESET, SUFFIX_REPLICATION_STATUS, + has_internal_suffix, insert_str, }; use rustfs_utils::{DEFAULT_SIP_HASH_KEY, get_env_usize, sip_hash}; #[cfg(test)] @@ -174,6 +175,14 @@ fn has_raw_status(err: &SdkError, status: u16) -> bool { err.raw_response().is_some_and(|r| r.status().as_u16() == status) } +fn metadata_requires_existing_target(op_type: ReplicationType, object_info: &ObjectInfo) -> bool { + op_type == ReplicationType::Metadata + && object_info + .user_defined + .get(AMZ_BUCKET_REPLICATION_STATUS) + .is_some_and(|status| status.eq_ignore_ascii_case(ReplicationStatusType::Replica.as_str())) +} + const METRIC_VERSION_IDENTITY_DRIFT_TOTAL: &str = "rustfs_replication_version_identity_drift_total"; /// Targets that already produced a version-identity-drift warning this @@ -3494,6 +3503,7 @@ async fn resolve_replicate_all_action( start_time, ssec_audit_required, } = ctx; + let require_existing_target = metadata_requires_existing_target(roi.op_type, &object_info); let replication_action; match head_object_for_worker(tgt_client.as_ref(), &tgt_client.bucket, object, roi.version_id.map(|v| v.to_string())).await { Ok(oi) => { @@ -3555,7 +3565,13 @@ async fn resolve_replicate_all_action( // Version-ID format mismatch: retry without versionId and compare ETags. match head_object_fallback(tgt_client, object).await { Ok(Some(oi)) => { - replication_action = if replication_etags_match(object_info.etag.as_deref(), oi.e_tag.as_deref()) { + let etags_match = replication_etags_match(object_info.etag.as_deref(), oi.e_tag.as_deref()); + if require_existing_target && !etags_match { + rinfo.error = Some("replica metadata target does not contain matching object data".to_string()); + rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs(); + return None; + } + replication_action = if etags_match { if ssec_audit_required && !settle_ssec_passthrough_evidence(&oi, tgt_client, bucket, object, rinfo).await { @@ -3568,6 +3584,11 @@ async fn resolve_replicate_all_action( }; } Ok(None) => { + if require_existing_target { + rinfo.error = Some("replica metadata target does not contain this object version".to_string()); + rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs(); + return None; + } replication_action = ReplicationAction::All; } Err(e2) => { @@ -3593,7 +3614,12 @@ async fn resolve_replicate_all_action( return None; } } - } else if e.as_service_error().is_some_and(|se| se.is_not_found()) { + } else if e.as_service_error().is_some_and(|se| se.is_not_found()) || has_raw_status(&e, 404) { + if require_existing_target { + rinfo.error = Some("replica metadata target does not contain this object version".to_string()); + rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs(); + return None; + } replication_action = ReplicationAction::All; } else { rinfo.error = Some(e.to_string()); @@ -3868,6 +3894,7 @@ async fn replicate_object_with_multipart(ctx: MultipartR actual_size, object_info.etag.clone().unwrap_or_default(), object_info.mod_time, + &put_opts.internal, ), ) .await @@ -3921,6 +3948,113 @@ mod tests { }) } + fn spawn_head_status_server(status: u16) -> (String, std::thread::JoinHandle<()>) { + use std::io::{Read, Write}; + + let listener = std::net::TcpListener::bind(("127.0.0.1", 0)).expect("test HTTP listener should bind"); + let endpoint = format!("http://{}", listener.local_addr().expect("test HTTP listener should have an address")); + let handle = std::thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("test HTTP client should connect"); + let mut request = [0_u8; 8192]; + let bytes_read = stream.read(&mut request).expect("test HTTP request should be read"); + assert!(bytes_read > 0, "test HTTP request should not be empty"); + assert!(request[..bytes_read].starts_with(b"HEAD "), "replication comparison must use HEAD"); + write!(stream, "HTTP/1.1 {status} Test\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") + .expect("test HTTP response should be written"); + }); + (endpoint, handle) + } + + #[tokio::test] + async fn replica_metadata_missing_target_stops_before_full_put() { + let (endpoint, server) = spawn_head_status_server(404); + let target = test_target_client(endpoint); + let roi = ReplicateObjectInfo { + bucket: "source".to_string(), + name: "object".to_string(), + version_id: Some(Uuid::new_v4()), + op_type: ReplicationType::Metadata, + // Normal metadata writes replace REPLICA with per-target PENDING + // before constructing the worker request. + replication_status: ReplicationStatusType::Pending, + ..Default::default() + }; + let object_info = ObjectInfo { + bucket: roi.bucket.clone(), + name: roi.name.clone(), + version_id: roi.version_id, + etag: Some("source-etag".to_string()), + user_defined: Arc::new(HashMap::from([( + AMZ_BUCKET_REPLICATION_STATUS.to_string(), + ReplicationStatusType::Replica.as_str().to_string(), + )])), + ..Default::default() + }; + let mut rinfo = replicate_all_target_info(&roi, &target); + + let action = resolve_replicate_all_action( + ReplicateAllActionContext { + roi: &roi, + tgt_client: &target, + bucket: &roi.bucket, + object: &roi.name, + start_time: OffsetDateTime::now_utc(), + ssec_audit_required: false, + }, + object_info, + &mut rinfo, + ) + .await; + + assert!(action.is_none(), "missing replica metadata targets must not reach the payload PUT path"); + assert_eq!(rinfo.replication_status, ReplicationStatusType::Failed); + assert_eq!( + rinfo.error.as_deref(), + Some("replica metadata target does not contain this object version") + ); + server.join().expect("test HTTP server should finish"); + } + + #[tokio::test] + async fn source_metadata_missing_target_rebuilds_object() { + let (endpoint, server) = spawn_head_status_server(404); + let target = test_target_client(endpoint); + let roi = ReplicateObjectInfo { + bucket: "source".to_string(), + name: "object".to_string(), + version_id: Some(Uuid::new_v4()), + op_type: ReplicationType::Metadata, + replication_status: ReplicationStatusType::Pending, + ..Default::default() + }; + let object_info = ObjectInfo { + bucket: roi.bucket.clone(), + name: roi.name.clone(), + version_id: roi.version_id, + etag: Some("source-etag".to_string()), + ..Default::default() + }; + let mut rinfo = replicate_all_target_info(&roi, &target); + + let action = resolve_replicate_all_action( + ReplicateAllActionContext { + roi: &roi, + tgt_client: &target, + bucket: &roi.bucket, + object: &roi.name, + start_time: OffsetDateTime::now_utc(), + ssec_audit_required: false, + }, + object_info, + &mut rinfo, + ) + .await; + + assert!(matches!(action, Some((ReplicationAction::All, _)))); + assert!(rinfo.error.is_none()); + server.join().expect("test HTTP server should finish"); + } + async fn register_test_target(target: &Arc) { ReplicationTargetStore::register_test_target(target).await; } diff --git a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs index 4fe89967d..50aa14cd7 100644 --- a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs @@ -472,6 +472,7 @@ pub(crate) fn replication_complete_multipart_options( actual_size: String, source_etag: String, source_mtime: Option, + 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() @@ -663,20 +672,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(), diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 5aa59b158..7e67997c0 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -2318,6 +2318,43 @@ 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. Only an absent version has no local state to compare; + // other read failures must leave the upload retryable rather than + // committing inbound metadata without the LWW check. + 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) => return Err(err), + } + } + for meta in parts_metadatas.iter_mut() { if meta.has_valid_erasure_geometry() { meta.size = fi.size; @@ -7055,6 +7092,146 @@ 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() + }; + + // Make the destination version unreadable on quorum while the + // staged upload remains intact. The commit barrier lets the old + // fail-open path move past the LWW read; restoring the metadata + // there proves it would otherwise commit the stale tags. + let mut damaged_metadata = Vec::new(); + for temp_dir in temp_dirs.iter().take(3) { + let path = temp_dir.path().join(bucket).join(object).join(STORAGE_FORMAT_FILE); + let original = tokio::fs::read(&path).await.expect("destination xl.meta should be readable"); + tokio::fs::write(&path, b"not an xl.meta") + .await + .expect("destination xl.meta should be corruptible"); + damaged_metadata.push((path, original)); + } + + let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::BeforeLockLost); + let first_set = set_disks.clone(); + let first_upload_id = upload_id.clone(); + let first_parts = parts.clone(); + let first_opts = complete_opts.clone(); + let mut first_completion = tokio::spawn(async move { + first_set + .complete_multipart_upload(bucket, object, &first_upload_id, first_parts, &first_opts) + .await + }); + + let first_result = tokio::select! { + result = &mut first_completion => result.expect("first completion task should finish"), + () = barrier.wait_until_paused() => { + for (path, original) in &damaged_metadata { + tokio::fs::write(path, original).await.expect("destination xl.meta should be restorable"); + } + barrier.release(); + first_completion.await.expect("released completion task should finish") + } + }; + for (path, original) in &damaged_metadata { + tokio::fs::write(path, original) + .await + .expect("destination xl.meta should be restored"); + } + drop(barrier); + + let first_error = first_result.expect_err("unreadable destination metadata must fail before multipart commit"); + assert!( + !(is_err_object_not_found(&first_error) || is_err_version_not_found(&first_error)), + "corrupt destination metadata must not be treated as an absent version: {first_error}" + ); + + 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() { diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 86f75eb8d..0730277d4 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -1881,6 +1881,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 { + 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, + existing: &HashMap, + 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, &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, @@ -2073,6 +2177,14 @@ impl SetDisks { user_defined.insert(key.clone(), value.clone()); } } + if replication_lww_applicable(opts) { + // Object Lock evaluation stamps category timestamps with this + // receiver's clock. Pin them back to the source-authored times + // before the first copy of a version is committed; the existing- + // version branch below may still replace them with newer local + // state. + merge_replication_metadata_lww(&mut user_defined, &HashMap::new(), opts); + } if expected_restore_operation_id.is_some() { rustfs_utils::http::metadata_compat::remove_str(&mut user_defined, SUFFIX_RESTORE_OPERATION_ID); } @@ -2562,6 +2674,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), @@ -8082,6 +8210,387 @@ 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, 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) -> 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, 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, 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 first_inbound_version_pins_source_authored_timestamp() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "lww-first-version-ts-pinned"; + let object = "object"; + let version_id = Uuid::new_v4().to_string(); + make_bucket(&disk_stores, bucket).await; + + let mut inbound = HashMap::new(); + inbound.insert(AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER.to_string(), "OFF".to_string()); + insert_str(&mut inbound, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP, T_OLD.to_string()); + let mut evaluated = inbound.clone(); + insert_str(&mut evaluated, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP, T_NEW.to_string()); + let opts = ObjectOptions { + replication_request: true, + replication_legalhold_timestamp: Some(parse_ts(T_OLD)), + eval_metadata: Some(evaluated), + ..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_LEGALHOLD_TIMESTAMP).as_deref(), + Some(T_OLD), + "the first copy must store the source timestamp, not the receiver evaluation time" + ); + } + + #[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; diff --git a/crates/replication/src/operation.rs b/crates/replication/src/operation.rs index a905ed992..ce566d9f8 100644 --- a/crates/replication/src/operation.rs +++ b/crates/replication/src/operation.rs @@ -88,6 +88,17 @@ impl MustReplicateOptions { return true; } + // A REPLICA version was delivered by a peer and carries no per-target + // internal status of its own; whether its local metadata edits flow + // back is the replication rule's ReplicaModifications decision + // (`ReplicationConfig::replicate` with `replica = true`, MinIO + // mustReplicate parity). Gating it on a COMPLETED target state would + // silently drop every replica-side tag / retention / legal-hold edit in + // an active-active topology (rustfs/backlog#1953). + if self.replication_status() == ReplicationStatusType::Replica { + return true; + } + get_internal_metadata(&self.meta, SUFFIX_REPLICATION_STATUS) .as_deref() .and_then(|statuses| { @@ -381,6 +392,13 @@ mod tests { assert!(options.metadata_target_is_eligible(arn)); assert!(!options.metadata_target_is_eligible("arn:rustfs:replication:missing")); + + // A replica-side metadata edit (active-active) has no per-target + // internal status; eligibility is left to the ReplicaModifications rule. + let replica = MustReplicateOptions::new(&HashMap::new(), String::new(), ReplicationType::Metadata, false) + .with_replication_status(ReplicationStatusType::Replica); + assert!(replica.metadata_target_is_eligible(arn)); + assert!(replica.metadata_target_is_eligible("arn:rustfs:replication:missing")); assert!( MustReplicateOptions::new(&HashMap::new(), String::new(), ReplicationType::Object, false) .metadata_target_is_eligible(arn) diff --git a/crates/utils/src/http/headers.rs b/crates/utils/src/http/headers.rs index b983218c0..380c59705 100644 --- a/crates/utils/src/http/headers.rs +++ b/crates/utils/src/http/headers.rs @@ -47,6 +47,8 @@ pub const AMZ_DELETE_MARKER: &str = "x-amz-delete-marker"; // S3 object tagging pub const AMZ_OBJECT_TAGGING: &str = "X-Amz-Tagging"; +/// Lowercase wire form of [`AMZ_OBJECT_TAGGING`] for `HeaderMap` insertion. +pub const AMZ_OBJECT_TAGGING_LOWER: &str = "x-amz-tagging"; pub const AMZ_TAG_COUNT: &str = "x-amz-tagging-count"; pub const AMZ_TAG_DIRECTIVE: &str = "X-Amz-Tagging-Directive"; diff --git a/rustfs/src/storage/options.rs b/rustfs/src/storage/options.rs index 701dbd38b..097cc30e0 100644 --- a/rustfs/src/storage/options.rs +++ b/rustfs/src/storage/options.rs @@ -554,10 +554,11 @@ fn apply_replication_timestamps_from_headers(headers: &HeaderMap, 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),