Compare commits

...

4 Commits

Author SHA1 Message Date
唐小鸭 74b1189016 fix(storage): reserve replication transport names at metadata ingest
Second review round: a client PUT of
x-amz-meta-x-rustfs-source-replication-tagging-timestamp materialized
the bare transport key as stored user metadata. The outbound
replication header builder forwards user metadata verbatim on a
server-authorized request, so the receiver would persist the
attacker-chosen value as trusted internal LWW state — and for a
tagless object nothing later overwrites it.

The ingest namespacing guard now reserves the whole
x-rustfs-source- / x-minio-source- families (the new timestamps and
their siblings: source-mtime/-etag/-version-id/-replication-request),
folding forged keys back under x-amz-meta-. Forged-ingress regression
covers both prefixes and a sibling.
2026-08-16 01:24:23 +08:00
唐小鸭 77a536093e fix(replication): load the stored tagging timestamp independently of remaining tags
Review: DeleteObjectTagging persists the tagging-timestamp internal key
but leaves the object tagless, and the outbound mapper only loaded the
key inside the user_tags-nonempty branch — the deletion's LWW timestamp
stayed at the epoch and the header was omitted, so the deletion could
never win conflict resolution on the replica. The stored key is now
loaded unconditionally; the mod_time fallback still applies only while
tags exist (MinIO parity), and a tagless object without the key keeps
the epoch default (no header). Deletion-path regression test added.
2026-08-15 18:51:39 +08:00
唐小鸭 8d191c7653 fix(replication): transport and persist LWW timestamps for tag, retention, and legal hold
Active-active conflict resolution for concurrent tag/retention/legal-hold
edits needs the source's per-category modification times on both sides of
the wire; the three AdvancedPutOptions timestamp fields were dead and the
headers were neither sent nor parsed.

- Emit x-{rustfs,minio}-source-replication-{tagging,retention,legalhold}-
  timestamp from PutObjectOptions::header(); names and RFC3339 values
  interoperate with MinIO (minio-go constants.go, object-api-options.go),
  pinned by a header_compat wire-name test.
- Default the three AdvancedPutOptions timestamps to UNIX_EPOCH and skip
  epoch values in header(), so "never modified" is not sent as a
  modification made now.
- Parse the headers only on authorized replication PUTs and multipart
  completes, expose them as Option<OffsetDateTime> on ObjectOptions, and
  persist them into the dual-prefix internal metadata keys so the
  outbound pass (replication_target_boundary) reads the source's
  timestamps instead of the mod_time fallback.
- Record the local tagging timestamp in the PutObjectTagging and
  DeleteObjectTagging eval metadata, mirroring the object-lock handlers;
  without it the sender only ever had the mod_time fallback to offer.

Receiver-side LWW comparison (keep newer stored category metadata over a
stale inbound copy) is left as a TODO at the parse site.
2026-08-15 10:40:54 +08:00
唐小鸭 faf735896b test(replication): pin missing LWW timestamp header transport
Red-light tests for the replication timestamp three-header contract:

- put_object_headers_carry_replication_timestamp_headers pins that
  PutObjectOptions::header() must emit the
  x-{rustfs,minio}-source-replication-{tagging,retention,legalhold}-timestamp
  headers when the internal timestamps are set (currently missing).
- test_put_opts_from_headers_gates_replication_timestamp_persistence_on_authorization
  and test_complete_multipart_opts_persist_replication_timestamps_when_authorized
  pin that an authorized replication PUT / multipart complete must persist
  the inbound timestamps into the internal metadata keys while unauthorized
  requests must not (currently never persisted).
- fake_s3_target journals the three timestamp headers per request
  (ReplicationTimestampHeaders on RequestRecord) so sender-side e2e
  assertions can observe what a real target receives; self-test included.
2026-08-15 10:08:25 +08:00
7 changed files with 448 additions and 17 deletions
+91 -1
View File
@@ -76,6 +76,18 @@ const SOURCE_MTIME_HEADERS: [&str; 2] = ["x-rustfs-source-mtime", "x-minio-sourc
const SOURCE_REPLICATION_REQUEST_HEADERS: [&str; 2] =
["x-rustfs-source-replication-request", "x-minio-source-replication-request"];
const SOURCE_ETAG_HEADERS: [&str; 2] = ["x-rustfs-source-etag", "x-minio-source-etag"];
const SOURCE_TAGGING_TIMESTAMP_HEADERS: [&str; 2] = [
"x-rustfs-source-replication-tagging-timestamp",
"x-minio-source-replication-tagging-timestamp",
];
const SOURCE_RETENTION_TIMESTAMP_HEADERS: [&str; 2] = [
"x-rustfs-source-replication-retention-timestamp",
"x-minio-source-replication-retention-timestamp",
];
const SOURCE_LEGALHOLD_TIMESTAMP_HEADERS: [&str; 2] = [
"x-rustfs-source-replication-legalhold-timestamp",
"x-minio-source-replication-legalhold-timestamp",
];
const RESERVED_BUCKET_PREFIXES: [&str; 3] = ["xn--", "sthree-", "amzn-s3-demo-"];
const RESERVED_BUCKET_SUFFIXES: [&str; 6] = ["-s3alias", "--ol-s3", ".mrap", "--x-s3", "--table-s3", "-an"];
@@ -118,6 +130,25 @@ pub enum FaultAction {
WrongEtag,
}
/// Replication LWW timestamp headers observed on a request, journaled so
/// sender-side tests can assert what a real target would receive.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ReplicationTimestampHeaders {
pub tagging: Option<String>,
pub retention: Option<String>,
pub legalhold: Option<String>,
}
impl ReplicationTimestampHeaders {
fn from_headers(headers: &HeaderMap) -> Self {
Self {
tagging: header_value(headers, &SOURCE_TAGGING_TIMESTAMP_HEADERS).map(bounded_journal_value),
retention: header_value(headers, &SOURCE_RETENTION_TIMESTAMP_HEADERS).map(bounded_journal_value),
legalhold: header_value(headers, &SOURCE_LEGALHOLD_TIMESTAMP_HEADERS).map(bounded_journal_value),
}
}
}
/// Credential-free request metadata retained for deterministic assertions.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RequestRecord {
@@ -131,6 +162,7 @@ pub struct RequestRecord {
pub part_number: Option<i32>,
pub content_length: Option<u64>,
pub consumed_bytes: Option<usize>,
pub replication_timestamps: ReplicationTimestampHeaders,
pub fault: Option<FaultAction>,
}
@@ -536,7 +568,15 @@ impl S3Access for FaultAccess {
.get(CONTENT_LENGTH)
.and_then(|value| value.to_str().ok())
.and_then(|value| value.parse().ok());
let fault = record_request(&self.control, operation, context.method().clone(), parsed, content_length);
let replication_timestamps = ReplicationTimestampHeaders::from_headers(context.headers());
let fault = record_request(
&self.control,
operation,
context.method().clone(),
parsed,
content_length,
replication_timestamps,
);
if let Some(RequestFault {
action: FaultAction::Status(status),
..
@@ -589,6 +629,7 @@ fn record_request(
method: Method,
parsed: ParsedRequest,
content_length: Option<u64>,
replication_timestamps: ReplicationTimestampHeaders,
) -> Option<RequestFault> {
let mut state = lock(control);
let action = parsed
@@ -613,6 +654,7 @@ fn record_request(
part_number: parsed.part_number,
content_length,
consumed_bytes: None,
replication_timestamps,
fault: action.clone(),
});
action.map(|action| RequestFault { sequence, action })
@@ -1699,6 +1741,52 @@ mod tests {
.await?)
}
#[tokio::test]
async fn journals_replication_timestamp_headers() -> Result<(), BoxError> {
let target = FakeS3Target::start().await?;
target.create_bucket("target-bucket");
let client = client(&target);
client
.put_object()
.bucket("target-bucket")
.key("plain")
.body(ByteStream::from_static(b"plain"))
.send()
.await?;
client
.put_object()
.bucket("target-bucket")
.key("stamped")
.body(ByteStream::from_static(b"stamped"))
.customize()
.map_request(move |mut request| {
let headers = request.headers_mut();
headers.insert("x-rustfs-source-replication-tagging-timestamp", "2026-01-02T03:04:05Z");
headers.insert("x-minio-source-replication-retention-timestamp", "2026-01-02T03:04:06Z");
headers.insert("x-rustfs-source-replication-legalhold-timestamp", "2026-01-02T03:04:07Z");
Ok::<_, std::convert::Infallible>(request)
})
.send()
.await?;
let requests = target.requests();
let plain = requests
.iter()
.find(|record| record.operation == Operation::PutObject && record.key.as_deref() == Some("plain"))
.expect("plain PUT must be journaled");
assert_eq!(plain.replication_timestamps, ReplicationTimestampHeaders::default());
let stamped = requests
.iter()
.find(|record| record.operation == Operation::PutObject && record.key.as_deref() == Some("stamped"))
.expect("stamped PUT must be journaled");
assert_eq!(stamped.replication_timestamps.tagging.as_deref(), Some("2026-01-02T03:04:05Z"));
assert_eq!(stamped.replication_timestamps.retention.as_deref(), Some("2026-01-02T03:04:06Z"));
assert_eq!(stamped.replication_timestamps.legalhold.as_deref(), Some("2026-01-02T03:04:07Z"));
Ok(())
}
macro_rules! assert_sdk_error {
($error:expr, $status:expr, $code:expr) => {{
let error = &$error;
@@ -2985,6 +3073,7 @@ mod tests {
part_number: None,
},
Some(0),
ReplicationTimestampHeaders::default(),
);
}
let records = lock(&control).requests.clone();
@@ -3006,6 +3095,7 @@ mod tests {
part_number: None,
},
None,
ReplicationTimestampHeaders::default(),
);
{
let bounded_records = lock(&bounded_control);
+70 -4
View File
@@ -58,7 +58,9 @@ use rustfs_utils::http::{
};
use rustfs_utils::http::{
SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_CHECK,
SUFFIX_SOURCE_REPLICATION_REQUEST, SUFFIX_SOURCE_VERSION_ID, insert_header,
SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, SUFFIX_SOURCE_VERSION_ID,
insert_header,
};
use rustls_pki_types::pem::PemObject;
use serde::{Deserialize, Serialize};
@@ -1476,9 +1478,12 @@ impl Default for AdvancedPutOptions {
replication_status: ReplicationStatusType::Pending,
source_mtime: OffsetDateTime::now_utc(),
replication_request: false,
retention_timestamp: OffsetDateTime::now_utc(),
tagging_timestamp: OffsetDateTime::now_utc(),
legalhold_timestamp: OffsetDateTime::now_utc(),
// UNIX_EPOCH means "never modified": header() must not emit a
// timestamp header for it, otherwise a receiver would treat an
// unset category as a modification made right now.
retention_timestamp: OffsetDateTime::UNIX_EPOCH,
tagging_timestamp: OffsetDateTime::UNIX_EPOCH,
legalhold_timestamp: OffsetDateTime::UNIX_EPOCH,
replication_validity_check: false,
}
}
@@ -1675,6 +1680,16 @@ impl PutObjectOptions {
);
}
for (suffix, timestamp) in [
(SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, self.internal.tagging_timestamp),
(SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP, self.internal.retention_timestamp),
(SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, self.internal.legalhold_timestamp),
] {
if timestamp.unix_timestamp() != 0 {
insert_header(&mut header, suffix, timestamp.format(&Rfc3339).unwrap_or_default());
}
}
if self.internal.replication_request {
insert_header(&mut header, SUFFIX_SOURCE_REPLICATION_REQUEST, "true");
}
@@ -2842,6 +2857,57 @@ mod tests {
);
}
#[test]
fn put_object_headers_carry_replication_timestamp_headers() {
// MinIO receivers resolve concurrent tag/retention/legal-hold edits by
// last-writer-wins on these headers (object-api-options.go parses them
// as RFC3339); a replica without them loses every conflict resolution.
let mut opts = PutObjectOptions::default();
opts.internal.replication_request = true;
let tagging = OffsetDateTime::from_unix_timestamp(1_700_000_001).expect("valid timestamp");
let retention = OffsetDateTime::from_unix_timestamp(1_700_000_002).expect("valid timestamp");
let legalhold = OffsetDateTime::from_unix_timestamp(1_700_000_003).expect("valid timestamp");
opts.internal.tagging_timestamp = tagging;
opts.internal.retention_timestamp = retention;
opts.internal.legalhold_timestamp = legalhold;
let header = opts.header();
for (suffix, expected) in [
("source-replication-tagging-timestamp", tagging),
("source-replication-retention-timestamp", retention),
("source-replication-legalhold-timestamp", legalhold),
] {
assert_eq!(
rustfs_utils::http::get_header(&header, suffix).as_deref(),
Some(expected.format(&Rfc3339).expect("RFC3339 timestamp").as_str()),
"replication put requests must carry the {suffix} header"
);
}
}
#[test]
fn put_object_headers_omit_unset_replication_timestamps() {
// UNIX_EPOCH means "never modified on the source"; sending it would
// make the receiver treat an unset category as a fresh modification.
let mut opts = PutObjectOptions::default();
opts.internal.replication_request = true;
opts.internal.tagging_timestamp = OffsetDateTime::UNIX_EPOCH;
opts.internal.retention_timestamp = OffsetDateTime::UNIX_EPOCH;
opts.internal.legalhold_timestamp = OffsetDateTime::UNIX_EPOCH;
let header = opts.header();
for suffix in [
"source-replication-tagging-timestamp",
"source-replication-retention-timestamp",
"source-replication-legalhold-timestamp",
] {
assert!(
rustfs_utils::http::get_header(&header, suffix).is_none(),
"unset {suffix} must not be sent to replication targets"
);
}
}
#[tokio::test]
async fn get_remote_target_client_internal_rejects_loopback_endpoint() {
let sys = BucketTargetSys::default();
@@ -259,15 +259,23 @@ pub(crate) fn replication_put_object_options(sc: &str, object_info: &ObjectInfo)
if !tags.is_empty() {
put_options.user_tags = tags;
put_options.internal.tagging_timestamp =
if let Some(timestamp) = get_str(&object_info.user_defined, SUFFIX_TAGGING_TIMESTAMP) {
OffsetDateTime::parse(&timestamp, &Rfc3339)
.map_err(|err| Error::other(format!("Failed to parse tagging timestamp: {err}")))?
} else {
object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH)
};
}
}
// Load the stored tagging timestamp independently of whether any tags
// remain: DeleteObjectTagging leaves the object tagless but stamps this
// key, and the deletion's LWW timestamp must still reach the replica.
// With no stored key, fall back to mod_time only while tags exist
// (MinIO parity); a tagless object without the key was never tagged and
// keeps the epoch default (no header).
put_options.internal.tagging_timestamp = if let Some(timestamp) = get_str(&object_info.user_defined, SUFFIX_TAGGING_TIMESTAMP)
{
OffsetDateTime::parse(&timestamp, &Rfc3339)
.map_err(|err| Error::other(format!("Failed to parse tagging timestamp: {err}")))?
} else if !put_options.user_tags.is_empty() {
object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH)
} else {
OffsetDateTime::UNIX_EPOCH
};
let metadata = &*object_info.user_defined;
@@ -694,6 +702,44 @@ mod tests {
assert!(options.internal.replication_request);
}
/// DeleteObjectTagging leaves the object tagless but stamps the
/// tagging-timestamp internal key; the deletion's LWW timestamp must
/// still be loaded (and therefore sent) so the replica can order the
/// deletion against concurrent tag edits.
#[test]
fn replication_put_options_carry_tagging_timestamp_after_tag_deletion() {
let mut metadata = std::collections::HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, SUFFIX_TAGGING_TIMESTAMP, "2026-01-02T03:04:05Z".to_string());
let object_info = ObjectInfo {
user_defined: Arc::new(metadata),
user_tags: Arc::new(String::new()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
version_id: Some(Uuid::nil()),
..Default::default()
};
let (options, _) = replication_put_object_options("", &object_info).expect("build put options");
assert!(options.user_tags.is_empty());
assert_eq!(
options.internal.tagging_timestamp,
OffsetDateTime::parse("2026-01-02T03:04:05Z", &Rfc3339).expect("valid timestamp"),
"the stored tagging timestamp must load independently of remaining tags"
);
// A tagless object without the stored key was never tagged: the epoch
// default keeps the header unsent.
let untagged = ObjectInfo {
user_tags: Arc::new(String::new()),
mod_time: Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp")),
version_id: Some(Uuid::nil()),
..Default::default()
};
let (options, _) = replication_put_object_options("", &untagged).expect("build put options");
assert_eq!(options.internal.tagging_timestamp, OffsetDateTime::UNIX_EPOCH);
}
#[test]
fn replication_put_options_strip_encryption_metadata_from_plaintext_objects() {
use rustfs_utils::http::object_encryption_keys::{INTERNAL_ENCRYPTION_ORIGINAL_SIZE_HEADER, SSEC_ORIGINAL_SIZE_HEADER};
+6
View File
@@ -277,6 +277,12 @@ pub struct ObjectOptions {
/// fence avoids recursively acquiring the read lock behind a queued writer.
pub bucket_lifecycle_lock_fence: Option<NamespaceLockFence>,
pub replication_request: bool,
/// Source-cluster LWW timestamps carried by an authorized replication
/// request; None when the source never modified the category. Only the
/// replication-authorized options builders may set these.
pub replication_tagging_timestamp: Option<OffsetDateTime>,
pub replication_retention_timestamp: Option<OffsetDateTime>,
pub replication_legalhold_timestamp: Option<OffsetDateTime>,
/// Authorized SSE-C replication passthrough: the body is already
/// ciphertext, so the write path must not encrypt or compress it and
/// stores the restored encryption metadata verbatim. Only the
+29
View File
@@ -50,6 +50,14 @@ pub const SUFFIX_SOURCE_DELETEMARKER: &str = "source-deletemarker";
pub const SUFFIX_SOURCE_PROXY_REQUEST: &str = "source-proxy-request";
pub const SUFFIX_SOURCE_REPLICATION_REQUEST: &str = "source-replication-request";
pub const SUFFIX_SOURCE_REPLICATION_CHECK: &str = "source-replication-check";
// LWW timestamps for replicated tag/retention/legal-hold modifications. MinIO
// declares these with mixed case (internal/http/headers.go:
// X-Minio-Source-Replication-Tagging-Timestamp / -Retention-Timestamp /
// -LegalHold-Timestamp); HTTP header names compare case-insensitively, so the
// lowercase suffix forms interoperate. Values are RFC3339 on the wire.
pub const SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP: &str = "source-replication-tagging-timestamp";
pub const SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP: &str = "source-replication-retention-timestamp";
pub const SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP: &str = "source-replication-legalhold-timestamp";
pub const SUFFIX_REPLICATION_SSEC_CRC: &str = "replication-ssec-crc";
/// Returns true if the key is object-encryption metadata understood by RustFS or MinIO.
@@ -196,6 +204,27 @@ mod tests {
assert_eq!(get_object_encryption_original_size(&metadata).expect("valid size"), Some(42));
}
#[test]
fn replication_timestamp_headers_match_minio_wire_names() {
let mut headers = HeaderMap::new();
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, "2026-01-02T03:04:05Z");
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP, "2026-01-02T03:04:06Z");
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, "2026-01-02T03:04:07Z");
// The exact names MinIO's object-api-options.go reads (its Get()
// canonicalizes case, so a case-insensitive match is wire-equivalent).
for name in [
"X-Minio-Source-Replication-Tagging-Timestamp",
"X-Minio-Source-Replication-Retention-Timestamp",
"X-Minio-Source-Replication-LegalHold-Timestamp",
"x-rustfs-source-replication-tagging-timestamp",
"x-rustfs-source-replication-retention-timestamp",
"x-rustfs-source-replication-legalhold-timestamp",
] {
assert!(headers.contains_key(name), "replication timestamp header {name} must be written");
}
}
#[test]
fn test_get_header() {
let mut headers = HeaderMap::new();
+11 -1
View File
@@ -45,7 +45,7 @@ use rustfs_targets::EventName;
use rustfs_utils::http::headers::{
AMZ_OBJECT_LOCK_LEGAL_HOLD_LOWER, AMZ_OBJECT_LOCK_MODE_LOWER, AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE_LOWER,
};
use rustfs_utils::http::{SUFFIX_REPLICATION_STATUS, SUFFIX_REPLICATION_TIMESTAMP, insert_str};
use rustfs_utils::http::{SUFFIX_REPLICATION_STATUS, SUFFIX_REPLICATION_TIMESTAMP, SUFFIX_TAGGING_TIMESTAMP, insert_str};
use s3s::{S3, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, dto::*, s3_error};
use std::collections::HashMap;
use std::fmt::Debug;
@@ -461,6 +461,11 @@ impl S3 for FS {
let mut eval_metadata = HashMap::new();
insert_str(&mut eval_metadata, SUFFIX_REPLICATION_TIMESTAMP, jiff::Zoned::now().to_string());
insert_str(&mut eval_metadata, SUFFIX_REPLICATION_STATUS, dsc.pending_status().unwrap_or_default());
insert_str(
&mut eval_metadata,
SUFFIX_TAGGING_TIMESTAMP,
OffsetDateTime::now_utc().format(&Rfc3339).unwrap_or_default(),
);
opts.eval_metadata = Some(eval_metadata);
}
@@ -1645,6 +1650,11 @@ impl S3 for FS {
let mut eval_metadata = HashMap::new();
insert_str(&mut eval_metadata, SUFFIX_REPLICATION_TIMESTAMP, jiff::Zoned::now().to_string());
insert_str(&mut eval_metadata, SUFFIX_REPLICATION_STATUS, dsc.pending_status().unwrap_or_default());
insert_str(
&mut eval_metadata,
SUFFIX_TAGGING_TIMESTAMP,
OffsetDateTime::now_utc().format(&Rfc3339).unwrap_or_default(),
);
opts.eval_metadata = Some(eval_metadata);
}
+188 -4
View File
@@ -17,11 +17,13 @@ use crate::storage::storage_api::options_consumer::contract::{object::HTTPPrecon
use http::header::{IF_MATCH, IF_NONE_MATCH};
use http::{HeaderMap, HeaderValue};
use rustfs_utils::http::{
AMZ_BUCKET_REPLICATION_STATUS, SUFFIX_FORCE_DELETE, SUFFIX_REPLICATION_ACTUAL_OBJECT_SIZE, SUFFIX_REPLICATION_SSEC_CRC,
SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_SOURCE_VERSION_ID, get_header,
AMZ_BUCKET_REPLICATION_STATUS, SUFFIX_FORCE_DELETE, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP,
SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, SUFFIX_REPLICATION_ACTUAL_OBJECT_SIZE, SUFFIX_REPLICATION_SSEC_CRC,
SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP,
SUFFIX_SOURCE_REPLICATION_REQUEST, SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP,
SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, SUFFIX_SOURCE_VERSION_ID, SUFFIX_TAGGING_TIMESTAMP, get_header,
header_compat::{MINIO_ENCRYPTION_PREFIX, RUSTFS_ENCRYPTION_PREFIX},
insert_header_map,
insert_header_map, insert_str,
metadata_compat::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX},
};
use rustfs_utils::http::{
@@ -422,6 +424,9 @@ pub fn get_complete_multipart_upload_opts_with_replication_authorization(
preserve_etag,
..Default::default()
};
if replication_request {
apply_replication_timestamps_from_headers(headers, &mut opts);
}
apply_replica_status_from_headers(headers, &mut opts, replication_request_authorized);
fill_conditional_writes_opts_from_header(headers, &mut opts)?;
@@ -479,6 +484,7 @@ pub fn put_opts_from_headers_with_replication_authorization(
if let Some(crc) = get_header(headers, SUFFIX_REPLICATION_SSEC_CRC) {
insert_header_map(&mut opts.user_defined, SUFFIX_REPLICATION_SSEC_CRC, crc.into_owned());
}
apply_replication_timestamps_from_headers(headers, &mut opts);
}
Ok(opts)
}
@@ -504,6 +510,47 @@ fn replication_source_mtime(headers: &HeaderMap<HeaderValue>) -> Option<time::Of
}
}
/// Parses one replication LWW timestamp header. Invalid values are dropped
/// with a warning (same tolerance as [`replication_source_mtime`]) so a
/// malformed source header cannot wedge the replication queue.
fn replication_timestamp_header(headers: &HeaderMap<HeaderValue>, suffix: &str) -> Option<time::OffsetDateTime> {
let value = get_header(headers, suffix)?;
let value = value.trim();
match time::OffsetDateTime::parse(value, &time::format_description::well_known::Rfc3339) {
Ok(timestamp) => Some(timestamp),
Err(err) => {
tracing::warn!("Invalid {} value '{}' (replication request=true): {}", suffix, value, err);
None
}
}
}
/// Callers must gate on an authorized replication request: these headers are
/// trusted source-cluster state, not client input.
fn apply_replication_timestamps_from_headers(headers: &HeaderMap<HeaderValue>, opts: &mut ObjectOptions) {
opts.replication_tagging_timestamp = replication_timestamp_header(headers, SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP);
opts.replication_retention_timestamp = replication_timestamp_header(headers, SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP);
opts.replication_legalhold_timestamp = replication_timestamp_header(headers, SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP);
// 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.
for (timestamp, suffix) in [
(opts.replication_tagging_timestamp, SUFFIX_TAGGING_TIMESTAMP),
(opts.replication_retention_timestamp, SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP),
(opts.replication_legalhold_timestamp, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP),
] {
if let Some(timestamp) = timestamp
&& let Ok(value) = timestamp.format(&time::format_description::well_known::Rfc3339)
{
insert_str(&mut opts.user_defined, suffix, value);
}
}
}
fn apply_replica_status_from_headers(headers: &HeaderMap<HeaderValue>, opts: &mut ObjectOptions, authorized: bool) {
if !authorized {
return;
@@ -663,6 +710,13 @@ fn is_reserved_user_metadata_key(key: &str) -> bool {
|| starts_with_ignore_ascii_case(key, MINIO_INTERNAL_PREFIX)
|| starts_with_ignore_ascii_case(key, RUSTFS_ENCRYPTION_PREFIX)
|| starts_with_ignore_ascii_case(key, MINIO_ENCRYPTION_PREFIX)
// Replication transport names (source-replication timestamps,
// source-mtime/-etag/-version-id, ...). A bare stored key with one of
// these names is forwarded verbatim by the outbound replication
// header builder on a server-authorized request, so the receiver
// would persist attacker-chosen values as trusted internal LWW state.
|| starts_with_ignore_ascii_case(key, "x-rustfs-source-")
|| starts_with_ignore_ascii_case(key, "x-minio-source-")
}
fn stored_user_metadata_key(key: &str) -> String {
@@ -1599,6 +1653,136 @@ mod tests {
assert!(opts_invalid.mod_time.is_none());
}
/// A client PUT must not materialize the replication transport names as
/// bare stored user-metadata keys: the outbound replication header
/// builder forwards user metadata verbatim on a server-authorized
/// request, so a bare `x-rustfs-source-replication-*-timestamp` key would
/// deliver an attacker-chosen value into the replica's trusted internal
/// LWW state (for a tagless object nothing later overwrites it).
#[test]
fn test_replication_transport_names_cannot_be_forged_via_user_metadata() {
let mut headers = HeaderMap::new();
for name in [
"x-amz-meta-x-rustfs-source-replication-tagging-timestamp",
"x-amz-meta-x-minio-source-replication-legalhold-timestamp",
"x-rustfs-meta-x-rustfs-source-replication-retention-timestamp",
"x-amz-meta-x-rustfs-source-mtime",
] {
headers.insert(
http::header::HeaderName::from_static(name),
HeaderValue::from_static("2026-01-02T03:04:05Z"),
);
}
let metadata = extract_metadata(&headers);
for forged in [
"x-rustfs-source-replication-tagging-timestamp",
"x-minio-source-replication-legalhold-timestamp",
"x-rustfs-source-replication-retention-timestamp",
"x-rustfs-source-mtime",
] {
assert!(
metadata.get(forged).is_none(),
"{forged} must not be storable as a bare user-metadata key"
);
}
// The values survive, namespaced back under the user-metadata prefix.
assert_eq!(
metadata
.get("x-amz-meta-x-rustfs-source-replication-tagging-timestamp")
.map(String::as_str),
Some("2026-01-02T03:04:05Z")
);
}
#[test]
fn test_put_opts_from_headers_gates_replication_timestamp_persistence_on_authorization() {
// Sender-side LWW state (replication_target_boundary.rs) is read back
// from these internal metadata keys, so an authorized replication PUT
// must persist the inbound timestamp headers; an unauthorized client
// must not be able to forge them.
let mut headers = HeaderMap::new();
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_REQUEST, "true");
insert_header(&mut headers, "source-replication-tagging-timestamp", "2026-01-02T03:04:05Z");
insert_header(&mut headers, "source-replication-retention-timestamp", "2026-01-02T03:04:06Z");
insert_header(&mut headers, "source-replication-legalhold-timestamp", "2026-01-02T03:04:07Z");
let untrusted = put_opts_from_headers(&headers, HashMap::new()).expect("ordinary PUT options should be created");
for suffix in [
"tagging-timestamp",
"objectlock-retention-timestamp",
"objectlock-legalhold-timestamp",
] {
assert!(
rustfs_utils::http::get_str(&untrusted.user_defined, suffix).is_none(),
"unauthorized clients must not persist the {suffix} internal key"
);
}
assert!(untrusted.replication_tagging_timestamp.is_none());
assert!(untrusted.replication_retention_timestamp.is_none());
assert!(untrusted.replication_legalhold_timestamp.is_none());
let trusted = put_opts_from_headers_with_replication_authorization(&headers, HashMap::new(), true)
.expect("authorized replication request should parse");
let parse = |value: &str| {
time::OffsetDateTime::parse(value, &time::format_description::well_known::Rfc3339).expect("valid RFC3339")
};
assert_eq!(trusted.replication_tagging_timestamp, Some(parse("2026-01-02T03:04:05Z")));
assert_eq!(trusted.replication_retention_timestamp, Some(parse("2026-01-02T03:04:06Z")));
assert_eq!(trusted.replication_legalhold_timestamp, Some(parse("2026-01-02T03:04:07Z")));
for (suffix, expected) in [
("tagging-timestamp", "2026-01-02T03:04:05Z"),
("objectlock-retention-timestamp", "2026-01-02T03:04:06Z"),
("objectlock-legalhold-timestamp", "2026-01-02T03:04:07Z"),
] {
assert_eq!(
rustfs_utils::http::get_str(&trusted.user_defined, suffix).as_deref(),
Some(expected),
"authorized replication must persist the {suffix} internal key"
);
}
}
#[test]
fn test_complete_multipart_opts_persist_replication_timestamps_when_authorized() {
let mut headers = HeaderMap::new();
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_REQUEST, "true");
insert_header(&mut headers, "replication-actual-object-size", "1");
insert_header(&mut headers, "source-replication-tagging-timestamp", "2026-01-02T03:04:05Z");
insert_header(&mut headers, "source-replication-retention-timestamp", "2026-01-02T03:04:06Z");
insert_header(&mut headers, "source-replication-legalhold-timestamp", "2026-01-02T03:04:07Z");
let untrusted = get_complete_multipart_upload_opts(&headers).expect("ordinary multipart options should be created");
for suffix in [
"tagging-timestamp",
"objectlock-retention-timestamp",
"objectlock-legalhold-timestamp",
] {
assert!(
rustfs_utils::http::get_str(&untrusted.user_defined, suffix).is_none(),
"unauthorized multipart completes must not persist the {suffix} internal key"
);
}
let trusted = get_complete_multipart_upload_opts_with_replication_authorization(&headers, true)
.expect("authorized multipart complete should parse");
for (suffix, expected) in [
("tagging-timestamp", "2026-01-02T03:04:05Z"),
("objectlock-retention-timestamp", "2026-01-02T03:04:06Z"),
("objectlock-legalhold-timestamp", "2026-01-02T03:04:07Z"),
] {
assert_eq!(
rustfs_utils::http::get_str(&trusted.user_defined, suffix).as_deref(),
Some(expected),
"authorized multipart completes must persist the {suffix} internal key"
);
}
assert!(trusted.replication_tagging_timestamp.is_some());
assert!(trusted.replication_retention_timestamp.is_some());
assert!(trusted.replication_legalhold_timestamp.is_some());
}
#[test]
fn test_put_opts_from_headers_with_replica_status() {
let mut headers = HeaderMap::new();