feat(replication): proxy GET/HEAD/Tagging for unreplicated objects to replication targets (#6172)

* feat(replication): proxy GET/HEAD/Tagging for unreplicated objects to replication targets

Implements the MinIO active-active read-proxy protocol (P1-5 of the
replication compatibility review): when a GET/HEAD/GetObjectTagging/
PutObjectTagging/DeleteObjectTagging request fails locally with
not-found and the bucket has replication targets, the request is proxied
to the targets in rule order, mirroring bucket-replication.go
proxyGetToReplicationTarget/proxyHeadToRepTarget/proxyTaggingToRepTarget.

Protocol surface:
- Anti-loop: inbound {x-rustfs-,x-minio-}source-proxy-request is parsed
  into ObjectOptions (proxy_request + proxy_header_set, matching MinIO
  ProxyRequest/ProxyHeaderSet); a request carrying the marker with ANY
  value is never re-proxied. Outbound client proxy calls send the marker
  as "true"; replication worker convergence HEADs send it as "false" so
  a peer's proxy layer cannot answer a convergence check by proxying
  back to the source (which would fake Completed without a PUT).
- Target selection: new replication_proxy.rs get_proxy_targets — empty
  when the marker is set, versioning is suspended, or no replication
  config; otherwise filter_target_arns -> TargetClient lookup, skipping
  targets with proxying disabled.
- TargetClient gains head_object_for_proxy/get_object (streaming) and
  the three tagging calls. Proxy calls never send the replication-check
  SSE-C exemption header; customer SSE-C keys are forwarded verbatim so
  the target performs real decryption. Conditional (If-*) headers are
  not forwarded (MinIO parity); Range and part_number are, with
  parts_count/tag_count/storage_class/expiration passed through.
- Metrics: proxy counters now count only real client proxy traffic,
  MinIO-aligned (one total per proxied request, one failed when no
  target served it). The previous misattributed counters — replication
  worker HEAD/PUT (#2672) and local tagging operations (#2682) — are
  removed; ReplProxyMetric now maps the tagging counters instead of
  dropping them.

e2e (fake_s3_target extended with tagging + header journaling): proxied
GET body + outbound header contract (marker present, no
replication-check, SSE-C passthrough), HEAD, anti-loop 404 with zero
outbound requests, GetObjectTagging, and metric mapping unit tests.

Rolling note: proxying only activates for buckets with replication
targets; requests carrying the marker keep pre-upgrade behavior.

Refs rustfs/backlog#1675 (P1-5)

* fix(replication): fail SSE-C passthrough closed on targets that drop transport headers (#6178)

SSE-C ciphertext passthrough replicates via X-Rustfs-Replication-* transport
headers. A MinIO/generic-S3 target silently discards them, storing bare
ciphertext with no decryption material — yet the PUT succeeded, so the object
reported COMPLETED with a silently unreadable replica (backlog#1675 N2).

Fail-closed design:
- SsecPassthroughCapability {Unknown, Supported, Unsupported} cached in
  BucketTargetSys per target ARN with a recording timestamp. Entries reset
  whenever the target is rebuilt, edited, or removed (arn_remotes_map
  lifecycle) and expire after SSEC_PASSTHROUGH_CAPABILITY_TTL (10 minutes):
  an expired verdict in either direction is re-earned through the audit, so
  an Unsupported target recovers automatically after an upgrade (at most one
  wasted PUT+HEAD audit per bad target per TTL window) and a Supported
  verdict cannot outlive a backend swapped behind the same endpoint.
- Replication worker (replicate_object and replicate_all): fresh Unsupported
  targets never receive the PUT — the attempt fails immediately into the
  normal MRF retry channel with a "run ?replication-check to re-probe" hint.
  Unknown or expired verdicts are audited: after the PUT the worker HEADs
  the replica back through the replication-check channel (source version id
  mapped through resolve_read_api_version_id, so null-version objects audit
  correctly) and requires SSE-C evidence (the echoed customer-algorithm
  header); missing evidence records Unsupported and fails the attempt.
  Convergence HEADs are audited the same way, so a broken ciphertext replica
  from an earlier attempt can never launder itself into COMPLETED via an
  ETag match. The gate/evidence policy is pure (replication_target_boundary,
  staleness folded in as an input) for the M2 worker migration.
- replication-check grows an SsecPassthrough probe phase: a probe PUT
  carrying the live transport-header shape, HEAD-back for evidence, and a
  machine-readable Code BucketRemoteSsecPassthroughUnsupported on failure.
  The probe verdict is synced into the runtime capability cache. Unlike
  VersionFidelity, a failed SsecPassthrough phase does NOT fail the target
  overall — it is a capability limit, not a broken replication contract,
  and a plaintext-only deployment against such a target must not turn red.
- fake_s3_target: default mode now models a RustFS target (stores the
  transport headers, echoes SSE-C evidence); the new
  drop_unlisted_replication_headers mode models MinIO. The journal records
  whether a request carried transport headers.

Receiver-echo verification: the replication-check HEAD exemption only skips
SSE-C key validation; the response has always built sse-customer-algorithm
from stored metadata (rustfs/src/app/object_usecase.rs), so no receiver
change was needed — pinned end to end by the replication-check e2e against
a real RustFS target.

Rolling-upgrade constraint: RustFS targets older than the replication-check
HEAD exemption (#5898) answer the audit HEAD without SSE-C evidence (or fail
it outright), so SSE-C replication to such targets reports FAILED. This is
deliberate — FAILED-and-retryable beats a silently undecryptable replica —
and self-heals: once the target is upgraded, the next TTL expiry (or a
manual ?replication-check re-probe) re-audits and records Supported.
Plaintext and managed-SSE replication are unaffected. The capability cache
is per-node; each node audits independently.

Known limitations:
- The audit judges evidence from the echoed customer-algorithm header only.
  A hypothetical target that preserves that one header while dropping other
  transport headers (partial-drop) would pass the audit; no known target
  behaves this way — observed targets drop the whole unknown-header family.
- A mixed-version target cluster can flap the verdict between audits routed
  to different target nodes until the rollout completes; the TTL bounds how
  long each stale verdict persists.

New e2e (backlog#1675 C1 + N2, red-first): fail-closed against a
header-dropping fake (FAILED + no second PUT via the capability cache,
journal-asserted; red run showed the old COMPLETED), replication-check
reports the SsecPassthrough phase Code while the target stays OK overall,
SSE-C heal convergence after a real target outage, and SSE-C
existing-object resync landing a REPLICA readable with the customer key.
TTL expiry in both directions is pinned at the cache and gate seams.
This commit is contained in:
唐小鸭
2026-08-18 08:47:22 +08:00
committed by GitHub
parent 21c2fb42bb
commit daecb93139
22 changed files with 2769 additions and 189 deletions
+2 -2
View File
@@ -4,8 +4,8 @@ This module is the shared failure-injection boundary for replication end-to-end
`FakeS3Target::start()` creates the listener. Add target buckets with `create_bucket`, point a RustFS remote target at `address()`, use `FAKE_ACCESS_KEY` / `FAKE_SECRET_KEY`, then enqueue per-operation faults with `inject`. Faults for one operation are consumed in FIFO order and do not consume faults queued for another operation. A fault is consumed only after `s3s` verifies the full request signature, so anonymous, other-access-key, and bad-signature traffic cannot disturb a script.
Supported data operations are HeadBucket, GetBucketVersioning, PUT/GET/HEAD/DELETE Object, and create/upload/complete/abort multipart upload. `create_bucket` models general-purpose buckets in S3's shared global namespace; account-regional namespace buckets and their `-an` names are intentionally out of scope. Buckets are versioned: PUT creates a version, DELETE without `versionId` creates a delete marker, and DELETE with `versionId` removes exactly that version. Internal source version IDs must be UUIDs and are stored canonically. Source mtime is honored only for source-replication PUT/DELETE requests; absent or invalid values use receipt time, matching RustFS, while multipart completion always uses receipt time. Replicated versions are ordered newest-first by source mtime so late older versions and delete markers do not become current. Equal mtimes prefer objects over delete markers, then canonical UUID order; RustFS's internal FileMeta signature tie-break is intentionally out of scope because it is not part of the target S3 protocol. Multipart part numbers follow S3's `1..=10000` range, and every completed part except the final part must be at least 5 MiB.
Supported data operations are HeadBucket, GetBucketVersioning, PUT/GET/HEAD/DELETE Object, Get/Put/Delete ObjectTagging (tags live per version; Put replaces the whole set, Delete clears it), and create/upload/complete/abort multipart upload. `create_bucket` models general-purpose buckets in S3's shared global namespace; account-regional namespace buckets and their `-an` names are intentionally out of scope. Buckets are versioned: PUT creates a version, DELETE without `versionId` creates a delete marker, and DELETE with `versionId` removes exactly that version. Internal source version IDs must be UUIDs and are stored canonically. Source mtime is honored only for source-replication PUT/DELETE requests; absent or invalid values use receipt time, matching RustFS, while multipart completion always uses receipt time. Replicated versions are ordered newest-first by source mtime so late older versions and delete markers do not become current. Equal mtimes prefer objects over delete markers, then canonical UUID order; RustFS's internal FileMeta signature tie-break is intentionally out of scope because it is not part of the target S3 protocol. Multipart part numbers follow S3's `1..=10000` range, and every completed part except the final part must be at least 5 MiB.
Fault actions cover HTTP 401/403/503 responses, pre-dispatch delay, connection abort when a logical request-body threshold is reached, streaming slow drain, and a deliberately wrong response ETag (including multipart-complete XML). `requests()` returns the ordered, credential-free request journal for assertions.
Fault actions cover HTTP 401/403/503 responses, pre-dispatch delay, connection abort when a logical request-body threshold is reached, streaming slow drain, and a deliberately wrong response ETag (including multipart-complete XML). `requests()` returns the ordered, credential-free request journal for assertions. Each record also journals a `ProxyHeaderSnapshot` — the read-proxy anti-loop marker (`x-{rustfs,minio}-source-proxy-request`), the replication-check exemption header, and the client SSE-C header family (algorithm and key-MD5 values; for the key itself only its presence) — so proxy tests can pin the exact wire contract.
The listener is loopback-only. It admits at most 64 active connections and two concurrently buffered request bodies; authenticated multipart-complete XML collection and assembly take both body permits. Keep-alive is disabled, request-header reads are bounded to 30 seconds, a parsed request is bounded to 65 seconds, and the complete connection lifetime is bounded to 100 seconds. It retains at most 256 buckets, 4,096 journal entries, 4,096 scripted faults, 4,096 object versions, 256 multipart uploads, and 10,000 multipart parts. Retained identifiers are capped at 1 KiB, user metadata at 2 KiB, and content type at 1 KiB. A PUT or uploaded part is capped at 64 MiB; a completed multipart object and all stored object/part data are capped at 128 MiB. Body drain, body-permit waits, delay, and slow-drain execution are bounded to 30 seconds; each slow-drain slice delay must be below that bound.
+284 -5
View File
@@ -30,10 +30,12 @@ use s3s::access::{S3Access, S3AccessContext};
use s3s::auth::SimpleAuth;
use s3s::dto::{
AbortMultipartUploadInput, AbortMultipartUploadOutput, CompleteMultipartUploadInput, CompleteMultipartUploadOutput,
CreateMultipartUploadInput, CreateMultipartUploadOutput, DeleteMarkerEntry, DeleteObjectInput, DeleteObjectOutput, ETag,
GetBucketVersioningInput, GetBucketVersioningOutput, GetObjectInput, GetObjectOutput, HeadBucketInput, HeadBucketOutput,
CreateMultipartUploadInput, CreateMultipartUploadOutput, DeleteMarkerEntry, DeleteObjectInput, DeleteObjectOutput,
DeleteObjectTaggingInput, DeleteObjectTaggingOutput, ETag, GetBucketVersioningInput, GetBucketVersioningOutput,
GetObjectInput, GetObjectOutput, GetObjectTaggingInput, GetObjectTaggingOutput, HeadBucketInput, HeadBucketOutput,
HeadObjectInput, HeadObjectOutput, ListObjectVersionsInput, ListObjectVersionsOutput, ObjectVersionId, PutObjectInput,
PutObjectOutput, StreamingBlob, Timestamp, TimestampFormat, UploadPartInput, UploadPartOutput,
PutObjectOutput, PutObjectTaggingInput, PutObjectTaggingOutput, StreamingBlob, Tag, TagSet, Timestamp, TimestampFormat,
UploadPartInput, UploadPartOutput,
};
use s3s::service::{S3Service, S3ServiceBuilder};
use s3s::validation::{AwsNameValidation, NameValidation};
@@ -88,6 +90,13 @@ const SOURCE_LEGALHOLD_TIMESTAMP_HEADERS: [&str; 2] = [
"x-rustfs-source-replication-legalhold-timestamp",
"x-minio-source-replication-legalhold-timestamp",
];
/// Wire prefix of the SSE-C passthrough replication transport headers
/// (`X-Rustfs-Replication-*`). In the default mode the fake stores them like a
/// RustFS target and echoes SSE-C evidence back on HEAD/GET; with
/// [`FakeS3Target::drop_unlisted_replication_headers`] it models MinIO /
/// generic S3, which silently discard unknown x-* headers.
const REPLICATION_SSE_TRANSPORT_PREFIX: &str = "x-rustfs-replication-";
const REPLICATION_SSEC_ALGORITHM_TRANSPORT_HEADER: &str = "x-rustfs-replication-ssec-algorithm";
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"];
@@ -103,6 +112,9 @@ pub enum Operation {
GetObject,
HeadObject,
DeleteObject,
GetObjectTagging,
PutObjectTagging,
DeleteObjectTagging,
ListObjectVersions,
CreateMultipartUpload,
UploadPart,
@@ -149,6 +161,42 @@ impl ReplicationTimestampHeaders {
}
}
/// Read-proxy related headers observed on a request, journaled so proxy
/// tests can assert the exact wire contract: the anti-loop marker present,
/// the replication-check exemption absent, and the client SSE-C key family
/// forwarded verbatim. The SSE-C key value itself is never retained — only
/// its presence.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ProxyHeaderSnapshot {
pub source_proxy_request: Option<String>,
pub replication_check: Option<String>,
pub ssec_algorithm: Option<String>,
pub ssec_key_present: bool,
pub ssec_key_md5: Option<String>,
/// Whether the request carried any `X-Rustfs-Replication-*` SSE-C
/// passthrough transport header, so fail-closed tests can assert the
/// sender really shipped the material a dropping target discarded.
pub ssec_transport_present: bool,
}
impl ProxyHeaderSnapshot {
fn from_headers(headers: &HeaderMap) -> Self {
Self {
source_proxy_request: header_value(headers, &["x-rustfs-source-proxy-request", "x-minio-source-proxy-request"])
.map(bounded_journal_value),
replication_check: header_value(headers, &["x-rustfs-source-replication-check", "x-minio-source-replication-check"])
.map(bounded_journal_value),
ssec_algorithm: header_value(headers, &["x-amz-server-side-encryption-customer-algorithm"])
.map(bounded_journal_value),
ssec_key_present: headers.contains_key("x-amz-server-side-encryption-customer-key"),
ssec_key_md5: header_value(headers, &["x-amz-server-side-encryption-customer-key-md5"]).map(bounded_journal_value),
ssec_transport_present: headers
.keys()
.any(|name| name.as_str().starts_with(REPLICATION_SSE_TRANSPORT_PREFIX)),
}
}
}
/// Credential-free request metadata retained for deterministic assertions.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RequestRecord {
@@ -163,6 +211,7 @@ pub struct RequestRecord {
pub content_length: Option<u64>,
pub consumed_bytes: Option<usize>,
pub replication_timestamps: ReplicationTimestampHeaders,
pub proxy_headers: ProxyHeaderSnapshot,
pub fault: Option<FaultAction>,
}
@@ -178,6 +227,10 @@ struct ControlState {
struct StoreState {
assign_own_version_ids: bool,
assign_own_multipart_version_ids: bool,
/// MinIO-like mode: silently discard non-whitelisted replication
/// transport headers instead of storing them (see
/// [`REPLICATION_SSE_TRANSPORT_PREFIX`]).
drop_unlisted_replication_headers: bool,
buckets: HashMap<String, BucketState>,
uploads: HashMap<String, MultipartState>,
total_bytes: usize,
@@ -199,6 +252,12 @@ struct ObjectVersion {
delete_marker: bool,
content_type: Option<String>,
metadata: Option<HashMap<String, String>>,
/// Object tags as ordered key/value pairs (PutObjectTagging replaces the
/// whole set, DeleteObjectTagging clears it).
tags: Vec<(String, String)>,
/// SSE-C passthrough transport headers stored with the version (RustFS
/// target behavior); empty when the drop mode discarded them.
replication_sse_headers: Vec<(String, String)>,
}
#[derive(Clone)]
@@ -208,6 +267,7 @@ struct MultipartState {
version_id: String,
content_type: Option<String>,
metadata: Option<HashMap<String, String>>,
replication_sse_headers: Vec<(String, String)>,
parts: BTreeMap<i32, MultipartPart>,
}
@@ -428,6 +488,15 @@ impl FakeS3Target {
/// Mint own version ids for the multipart path only — models a target
/// that adopts PutObject version ids but not CreateMultipartUpload ones.
/// MinIO-like mode: silently drop every `X-Rustfs-Replication-*` SSE-C
/// passthrough transport header instead of storing it. The default (off)
/// models a RustFS target, which preserves the headers and echoes SSE-C
/// evidence (`x-amz-server-side-encryption-customer-algorithm`) on
/// HEAD/GET of the replica.
pub fn drop_unlisted_replication_headers(&self, enabled: bool) {
lock(&self.backend.store).drop_unlisted_replication_headers = enabled;
}
pub fn assign_own_multipart_version_ids(&self, enabled: bool) {
lock(&self.backend.store).assign_own_multipart_version_ids = enabled;
}
@@ -569,6 +638,7 @@ impl S3Access for FaultAccess {
.and_then(|value| value.to_str().ok())
.and_then(|value| value.parse().ok());
let replication_timestamps = ReplicationTimestampHeaders::from_headers(context.headers());
let proxy_headers = ProxyHeaderSnapshot::from_headers(context.headers());
let fault = record_request(
&self.control,
operation,
@@ -576,6 +646,7 @@ impl S3Access for FaultAccess {
parsed,
content_length,
replication_timestamps,
proxy_headers,
);
if let Some(RequestFault {
action: FaultAction::Status(status),
@@ -615,6 +686,9 @@ fn operation_from_s3_name(name: &str) -> Operation {
"GetObject" => Operation::GetObject,
"HeadObject" => Operation::HeadObject,
"DeleteObject" => Operation::DeleteObject,
"GetObjectTagging" => Operation::GetObjectTagging,
"PutObjectTagging" => Operation::PutObjectTagging,
"DeleteObjectTagging" => Operation::DeleteObjectTagging,
"CreateMultipartUpload" => Operation::CreateMultipartUpload,
"UploadPart" => Operation::UploadPart,
"CompleteMultipartUpload" => Operation::CompleteMultipartUpload,
@@ -630,6 +704,7 @@ fn record_request(
parsed: ParsedRequest,
content_length: Option<u64>,
replication_timestamps: ReplicationTimestampHeaders,
proxy_headers: ProxyHeaderSnapshot,
) -> Option<RequestFault> {
let mut state = lock(control);
let action = parsed
@@ -655,6 +730,7 @@ fn record_request(
content_length,
consumed_bytes: None,
replication_timestamps,
proxy_headers,
fault: action.clone(),
});
action.map(|action| RequestFault { sequence, action })
@@ -721,6 +797,15 @@ fn parse_request(method: &Method, uri: &Uri) -> ParsedRequest {
(&Method::POST, true) if query.contains_key("uploads") => Operation::CreateMultipartUpload,
(&Method::POST, true) if upload_id.is_some() => Operation::CompleteMultipartUpload,
(&Method::DELETE, true) if upload_id.is_some() => Operation::AbortMultipartUpload,
(&Method::GET, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::GetObjectTagging
}
(&Method::PUT, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::PutObjectTagging
}
(&Method::DELETE, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::DeleteObjectTagging
}
// A replication PUT addresses the source version via `?versionId=`.
(&Method::PUT, true) if only_query_keys(&["versionId"]) => Operation::PutObject,
(&Method::GET, true) if only_query_keys(&["versionId"]) => Operation::GetObject,
@@ -788,6 +873,29 @@ fn new_version_id(headers: &HeaderMap, assign_own: bool) -> S3Result<String> {
Ok(version_id.to_string())
}
/// Capture the SSE-C passthrough transport headers a replication PUT carried.
/// Returns an empty set in the MinIO-like drop mode.
fn captured_replication_sse_headers(headers: &HeaderMap, drop_unlisted: bool) -> Vec<(String, String)> {
if drop_unlisted {
return Vec::new();
}
headers
.iter()
.filter(|(name, _)| name.as_str().starts_with(REPLICATION_SSE_TRANSPORT_PREFIX))
.filter_map(|(name, value)| Some((name.as_str().to_string(), value.to_str().ok()?.to_string())))
.collect()
}
/// SSE-C evidence a RustFS-like target echoes for a stored passthrough
/// replica: the customer algorithm restored from the transport headers.
fn stored_sse_customer_algorithm(version: &ObjectVersion) -> Option<String> {
version
.replication_sse_headers
.iter()
.find(|(name, _)| name == REPLICATION_SSEC_ALGORITHM_TRANSPORT_HEADER)
.map(|(_, value)| value.clone())
}
fn source_etag(headers: &HeaderMap) -> S3Result<Option<String>> {
header_value(headers, &SOURCE_ETAG_HEADERS)
.map(|value| validate_retained_identifier(value, "source ETag").map(|value| normalize_etag(&value)))
@@ -1135,6 +1243,33 @@ fn find_version(state: &StoreState, bucket: &str, key: &str, version_id: Option<
Ok(version.clone())
}
/// Replace (or clear, with an empty vec) the tag set of the addressed
/// version, returning its version id. Mirrors `find_version` addressing:
/// explicit version id or the latest version, delete markers rejected.
fn set_version_tags(
state: &mut StoreState,
bucket: &str,
key: &str,
version_id: Option<&str>,
tags: Vec<(String, String)>,
) -> S3Result<String> {
// Resolve first (immutable) so the error paths match find_version.
let resolved = find_version(state, bucket, key, version_id)?.version_id;
let versions = state
.buckets
.get_mut(bucket)
.expect("bucket existence checked by find_version")
.objects
.get_mut(key)
.expect("key existence checked by find_version");
let version = versions
.iter_mut()
.find(|version| version.version_id == resolved)
.expect("version existence checked by find_version");
version.tags = tags;
Ok(resolved)
}
#[async_trait]
impl S3 for FakeBackend {
async fn head_bucket(&self, req: S3Request<HeadBucketInput>) -> S3Result<S3Response<HeadBucketOutput>> {
@@ -1231,7 +1366,10 @@ impl S3 for FakeBackend {
let input = req.input;
let body = collect_stream(input.body, input.content_length, fault.as_ref(), &self.control).await?;
validate_stored_metadata(&input.content_type, &input.metadata)?;
let assign_own = lock(&self.store).assign_own_version_ids;
let (assign_own, drop_unlisted) = {
let state = lock(&self.store);
(state.assign_own_version_ids, state.drop_unlisted_replication_headers)
};
let version_id = new_version_id(&headers, assign_own)?;
let e_tag = match source_etag(&headers)? {
Some(value) => value,
@@ -1248,6 +1386,8 @@ impl S3 for FakeBackend {
delete_marker: false,
content_type: input.content_type,
metadata: input.metadata,
tags: Vec::new(),
replication_sse_headers: captured_replication_sse_headers(&headers, drop_unlisted),
};
upsert_version(&mut lock(&self.store), &input.bucket, input.key, version)?;
Ok(apply_response_fault(
@@ -1268,6 +1408,7 @@ impl S3 for FakeBackend {
let state = lock(&self.store);
find_version(&state, &input.bucket, &input.key, input.version_id.as_deref())?
};
let sse_customer_algorithm = stored_sse_customer_algorithm(&version);
Ok(apply_response_fault(
S3Response::new(GetObjectOutput {
body: Some(StreamingBlob::new(Body::from(version.body.clone()))),
@@ -1277,6 +1418,7 @@ impl S3 for FakeBackend {
e_tag: Some(ETag::Strong(version.e_tag)),
last_modified: Some(version.last_modified.clone()),
version_id: Some(version.version_id),
sse_customer_algorithm,
..Default::default()
}),
fault.as_ref(),
@@ -1291,6 +1433,7 @@ impl S3 for FakeBackend {
let state = lock(&self.store);
find_version(&state, &input.bucket, &input.key, input.version_id.as_deref())?
};
let sse_customer_algorithm = stored_sse_customer_algorithm(&version);
Ok(apply_response_fault(
S3Response::new(HeadObjectOutput {
content_length: Some(version.body.len() as i64),
@@ -1299,12 +1442,79 @@ impl S3 for FakeBackend {
e_tag: Some(ETag::Strong(version.e_tag)),
last_modified: Some(version.last_modified.clone()),
version_id: Some(version.version_id),
sse_customer_algorithm,
..Default::default()
}),
fault.as_ref(),
))
}
async fn get_object_tagging(&self, req: S3Request<GetObjectTaggingInput>) -> S3Result<S3Response<GetObjectTaggingOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let version = {
let state = lock(&self.store);
find_version(&state, &input.bucket, &input.key, input.version_id.as_deref())?
};
let tag_set: TagSet = version
.tags
.into_iter()
.map(|(key, value)| Tag {
key: Some(key),
value: Some(value),
})
.collect();
Ok(apply_response_fault(
S3Response::new(GetObjectTaggingOutput {
tag_set,
version_id: Some(ObjectVersionId::from(version.version_id)),
}),
fault.as_ref(),
))
}
async fn put_object_tagging(&self, req: S3Request<PutObjectTaggingInput>) -> S3Result<S3Response<PutObjectTaggingOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let tags = input
.tagging
.tag_set
.into_iter()
.map(|tag| (tag.key.unwrap_or_default(), tag.value.unwrap_or_default()))
.collect();
let version_id = {
let mut state = lock(&self.store);
set_version_tags(&mut state, &input.bucket, &input.key, input.version_id.as_deref(), tags)?
};
Ok(apply_response_fault(
S3Response::new(PutObjectTaggingOutput {
version_id: Some(ObjectVersionId::from(version_id)),
}),
fault.as_ref(),
))
}
async fn delete_object_tagging(
&self,
req: S3Request<DeleteObjectTaggingInput>,
) -> S3Result<S3Response<DeleteObjectTaggingOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let version_id = {
let mut state = lock(&self.store);
set_version_tags(&mut state, &input.bucket, &input.key, input.version_id.as_deref(), Vec::new())?
};
Ok(apply_response_fault(
S3Response::new(DeleteObjectTaggingOutput {
version_id: Some(ObjectVersionId::from(version_id)),
}),
fault.as_ref(),
))
}
async fn delete_object(&self, req: S3Request<DeleteObjectInput>) -> S3Result<S3Response<DeleteObjectOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
@@ -1381,6 +1591,8 @@ impl S3 for FakeBackend {
delete_marker: true,
content_type: None,
metadata: None,
tags: Vec::new(),
replication_sse_headers: Vec::new(),
},
)?;
Ok(apply_response_fault(
@@ -1408,9 +1620,10 @@ impl S3 for FakeBackend {
ensure_upload_budget(&state)?;
validate_stored_metadata(&input.content_type, &input.metadata)?;
let upload_id = Uuid::new_v4().to_string();
// Read the flag before the mutable borrow of `state.uploads` below
// Read the flags before the mutable borrow of `state.uploads` below
// (and never re-lock the store: the mutex is not reentrant).
let mint_own = state.assign_own_version_ids || state.assign_own_multipart_version_ids;
let drop_unlisted = state.drop_unlisted_replication_headers;
let version_id = new_version_id(&headers, mint_own)?;
state.uploads.insert(
upload_id.clone(),
@@ -1420,6 +1633,7 @@ impl S3 for FakeBackend {
version_id,
content_type: input.content_type,
metadata: input.metadata,
replication_sse_headers: captured_replication_sse_headers(&headers, drop_unlisted),
parts: BTreeMap::new(),
},
);
@@ -1557,6 +1771,7 @@ impl S3 for FakeBackend {
version_id: upload.version_id.clone(),
content_type: upload.content_type.clone(),
metadata: upload.metadata.clone(),
replication_sse_headers: upload.replication_sse_headers.clone(),
parts: BTreeMap::new(),
},
selected,
@@ -1583,6 +1798,8 @@ impl S3 for FakeBackend {
delete_marker: false,
content_type: upload.content_type,
metadata: upload.metadata,
tags: Vec::new(),
replication_sse_headers: upload.replication_sse_headers,
};
let mut state = lock(&self.store);
let current = state
@@ -1787,6 +2004,65 @@ mod tests {
Ok(())
}
/// Default mode is RustFS-like: SSE-C passthrough transport headers are
/// stored and the customer algorithm is echoed on HEAD/GET. Drop mode is
/// MinIO-like: the headers are silently discarded, so no evidence comes
/// back — the exact difference the N2 fail-closed audit keys on. Both
/// modes journal that the sender shipped the transport headers.
#[tokio::test]
async fn ssec_passthrough_headers_echo_and_drop_modes() -> Result<(), BoxError> {
let target = FakeS3Target::start().await?;
target.create_bucket("target-bucket");
let client = client(&target);
let put_with_transport_headers = |key: &'static str| {
client
.put_object()
.bucket("target-bucket")
.key(key)
.body(ByteStream::from_static(b"ciphertext"))
.customize()
.map_request(move |mut request| {
let headers = request.headers_mut();
headers.insert("x-rustfs-replication-ssec-algorithm", "AES256");
headers.insert("x-rustfs-replication-ssec-key-md5", "AAAAAAAAAAAAAAAAAAAAAA==");
Ok::<_, std::convert::Infallible>(request)
})
.send()
};
put_with_transport_headers("kept").await?;
let head = client.head_object().bucket("target-bucket").key("kept").send().await?;
assert_eq!(head.sse_customer_algorithm(), Some("AES256"));
let get = client.get_object().bucket("target-bucket").key("kept").send().await?;
assert_eq!(get.sse_customer_algorithm(), Some("AES256"));
target.drop_unlisted_replication_headers(true);
put_with_transport_headers("dropped").await?;
let head = client.head_object().bucket("target-bucket").key("dropped").send().await?;
assert_eq!(head.sse_customer_algorithm(), None, "drop mode must discard SSE-C evidence");
let requests = target.requests();
for key in ["kept", "dropped"] {
let record = requests
.iter()
.find(|record| record.operation == Operation::PutObject && record.key.as_deref() == Some(key))
.expect("PUT must be journaled");
assert!(
record.proxy_headers.ssec_transport_present,
"the journal must prove the sender shipped the transport headers for {key}"
);
}
let plain_head = requests
.iter()
.find(|record| record.operation == Operation::HeadObject)
.expect("HEAD must be journaled");
assert!(!plain_head.proxy_headers.ssec_transport_present);
target.shutdown().await;
Ok(())
}
macro_rules! assert_sdk_error {
($error:expr, $status:expr, $code:expr) => {{
let error = &$error;
@@ -3052,6 +3328,7 @@ mod tests {
version_id: index.to_string(),
content_type: None,
metadata: None,
replication_sse_headers: Vec::new(),
parts: BTreeMap::new(),
},
);
@@ -3074,6 +3351,7 @@ mod tests {
},
Some(0),
ReplicationTimestampHeaders::default(),
ProxyHeaderSnapshot::default(),
);
}
let records = lock(&control).requests.clone();
@@ -3096,6 +3374,7 @@ mod tests {
},
None,
ReplicationTimestampHeaders::default(),
ProxyHeaderSnapshot::default(),
);
{
let bounded_records = lock(&bounded_control);
@@ -2610,17 +2610,20 @@ async fn test_replication_check_succeeds_with_remote_target() -> Result<(), Box<
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await?;
assert_eq!(payload["Status"], "OK");
assert_eq!(payload["Status"], "OK", "{payload}");
assert_eq!(payload["ActiveMutation"], true);
assert_eq!(payload["Targets"].as_array().map(Vec::len), Some(1));
assert_eq!(payload["Targets"][0]["Status"], "OK");
assert_eq!(payload["Targets"][0]["Phases"]["Put"]["Status"], "OK");
assert_eq!(payload["Targets"][0]["Status"], "OK", "{payload}");
assert_eq!(payload["Targets"][0]["Phases"]["Put"]["Status"], "OK", "{payload}");
// A RustFS target adopts the source version id, so the P1-19
// version-identity probe passes.
assert_eq!(payload["Targets"][0]["Phases"]["VersionFidelity"]["Status"], "OK");
assert_eq!(payload["Targets"][0]["Phases"]["DeleteMarker"]["Status"], "OK");
assert_eq!(payload["Targets"][0]["Phases"]["VersionDelete"]["Status"], "OK");
assert_eq!(payload["Targets"][0]["Phases"]["Cleanup"]["Status"], "OK");
assert_eq!(payload["Targets"][0]["Phases"]["VersionFidelity"]["Status"], "OK", "{payload}");
// A RustFS target preserves the SSE-C passthrough transport headers and
// echoes the customer algorithm on the replication-check HEAD (N2).
assert_eq!(payload["Targets"][0]["Phases"]["SsecPassthrough"]["Status"], "OK", "{payload}");
assert_eq!(payload["Targets"][0]["Phases"]["DeleteMarker"]["Status"], "OK", "{payload}");
assert_eq!(payload["Targets"][0]["Phases"]["VersionDelete"]["Status"], "OK", "{payload}");
assert_eq!(payload["Targets"][0]["Phases"]["Cleanup"]["Status"], "OK", "{payload}");
let target_client = target_env.create_s3_client();
let versions = target_client
@@ -4649,6 +4652,410 @@ async fn test_bucket_replication_sse_c_multipart_passthrough() -> TestResult {
Ok(())
}
/// N2 (backlog#1675 P1-22): SSE-C passthrough replication to a target that
/// silently drops the `X-Rustfs-Replication-*` transport headers (MinIO-like
/// behavior, modeled by the fake target's drop mode) used to report COMPLETED
/// while the replica had irrecoverably lost its decryption material — the red
/// light this test was born failing on. Fail-closed contract now under test:
/// the first attempt PUTs, HEAD-backs the replica, finds no SSE-C evidence,
/// records the target Unsupported and reports FAILED; a second SSE-C object
/// fails without any PUT reaching the target (capability cache, proven from
/// the target journal); plaintext objects still replicate COMPLETED.
#[tokio::test]
#[serial]
async fn test_ssec_replication_fails_closed_when_target_drops_passthrough_headers() -> TestResult {
init_logging();
let target = FakeS3Target::start().await?;
let target_bucket = "ssec-drop-dst";
target.create_bucket(target_bucket);
target.drop_unlisted_replication_headers(true);
let mut source_env = RustFSTestEnvironment::new().await?;
let mut env_vars = replication_fast_env();
env_vars.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
env_vars.extend_from_slice(&[("NO_PROXY", "127.0.0.1,localhost"), ("HTTP_PROXY", ""), ("HTTPS_PROXY", "")]);
source_env.start_rustfs_server_with_env(vec![], &env_vars).await?;
let source_bucket = "ssec-drop-src";
let source_client = source_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
let target_arn = set_replication_target_with_options(
&source_env,
source_bucket,
ReplicationTargetOptions {
endpoint: &target.address(),
access_key: FAKE_ACCESS_KEY,
secret_key: FAKE_SECRET_KEY,
target_bucket,
secure: false,
skip_tls_verify: false,
ca_cert_pem: None,
},
)
.await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let customer_key = BASE64_STANDARD.encode(REPL17_SSEC_KEY);
let customer_key_md5 = sse_customer_key_md5_base64(REPL17_SSEC_KEY);
let put_ssec = |key: &'static str| {
source_client
.put_object()
.bucket(source_bucket)
.key(key)
.body(ByteStream::from_static(b"ssec fail-closed payload"))
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
};
// First SSE-C object: the audit must catch the dropped material.
put_ssec("ssec-first.txt").await?;
wait_for_source_replication_status(&source_client, source_bucket, "ssec-first.txt", "FAILED", true).await?;
let requests = target.take_requests();
let first_put = requests
.iter()
.find(|record| record.operation == FakeTargetOperation::PutObject && record.key.as_deref() == Some("ssec-first.txt"))
.ok_or("the first SSE-C object must have been PUT (capability was Unknown)")?;
assert!(
first_put.proxy_headers.ssec_transport_present,
"the replication PUT must have shipped the SSE-C transport headers the target then dropped"
);
assert!(
requests.iter().any(|record| {
record.operation == FakeTargetOperation::HeadObject
&& record.key.as_deref() == Some("ssec-first.txt")
&& record.sequence > first_put.sequence
&& record.proxy_headers.replication_check.as_deref() == Some("true")
}),
"the post-PUT HEAD-back audit must have run through the replication-check channel; journal: {requests:?}"
);
// Second SSE-C object: the cached Unsupported verdict fails it closed
// before any PUT — including MRF retries of the first object.
put_ssec("ssec-second.txt").await?;
wait_for_source_replication_status(&source_client, source_bucket, "ssec-second.txt", "FAILED", true).await?;
assert!(
!target.requests().iter().any(|record| {
record.operation == FakeTargetOperation::PutObject
&& record.key.as_deref() != Some("plain-control.txt")
&& record.proxy_headers.ssec_transport_present
}),
"no further SSE-C ciphertext may reach a target recorded Unsupported; journal: {:?}",
target.requests()
);
// The gate is scoped to SSE-C: plaintext replication keeps working.
source_client
.put_object()
.bucket(source_bucket)
.key("plain-control.txt")
.body(ByteStream::from_static(b"plaintext control payload"))
.send()
.await?;
wait_for_source_replication_status(&source_client, source_bucket, "plain-control.txt", "COMPLETED", false).await?;
assert!(target.has_object(target_bucket, "plain-control.txt"));
target.shutdown().await;
Ok(())
}
/// N2 (backlog#1675 P1-22): the admin replication-check must expose the same
/// verdict operators would otherwise only learn from failing SSE-C objects —
/// an SsecPassthrough probe phase that fails with the machine-readable
/// `BucketRemoteSsecPassthroughUnsupported` code against a header-dropping
/// target, with no probe residue left behind. The target's overall status
/// stays OK: unlike version-identity drift, dropped passthrough headers are
/// a capability limit, and a plaintext-only deployment against a MinIO-like
/// target must not turn red.
#[tokio::test]
#[serial]
async fn test_replication_check_flags_ssec_passthrough_dropping_target() -> TestResult {
init_logging();
let target = FakeS3Target::start().await?;
let target_bucket = "ssec-check-dst";
target.create_bucket(target_bucket);
target.drop_unlisted_replication_headers(true);
let mut source_env = RustFSTestEnvironment::new().await?;
let mut env_vars = replication_fast_env();
env_vars.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
env_vars.extend_from_slice(&[("NO_PROXY", "127.0.0.1,localhost"), ("HTTP_PROXY", ""), ("HTTPS_PROXY", "")]);
source_env.start_rustfs_server_with_env(vec![], &env_vars).await?;
let source_bucket = "ssec-check-src";
let source_client = source_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
let target_arn = set_replication_target_with_options(
&source_env,
source_bucket,
ReplicationTargetOptions {
endpoint: &target.address(),
access_key: FAKE_ACCESS_KEY,
secret_key: FAKE_SECRET_KEY,
target_bucket,
secure: false,
skip_tls_verify: false,
ca_cert_pem: None,
},
)
.await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let response = run_replication_check(&source_env, source_bucket).await?;
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await?;
assert_eq!(
payload["Status"], "OK",
"a capability-only SSE-C failure must not fail the check overall: {payload}"
);
let target_report = &payload["Targets"][0];
assert_eq!(target_report["Status"], "OK", "{payload}");
let ssec = &target_report["Phases"]["SsecPassthrough"];
assert_eq!(ssec["Status"], "FAILED", "SsecPassthrough phase must fail: {payload}");
assert_eq!(
ssec["Code"], "BucketRemoteSsecPassthroughUnsupported",
"the failure must carry the machine-readable code: {payload}"
);
// Basic replication of plaintext objects works on this target: every other
// phase passes, so the code is the discriminator operators branch on.
assert_eq!(target_report["Phases"]["Put"]["Status"], "OK", "{payload}");
assert_eq!(target_report["Phases"]["VersionFidelity"]["Status"], "OK", "{payload}");
assert_eq!(target_report["Phases"]["DeleteMarker"]["Status"], "OK", "{payload}");
assert_eq!(target_report["Phases"]["VersionDelete"]["Status"], "OK", "{payload}");
assert_eq!(target_report["Phases"]["Cleanup"]["Status"], "OK", "{payload}");
// The SSE-C probe PUT must have shipped the real transport header names —
// a mangled or missing header set would fail the phase for the wrong
// reason and mask a working target.
let requests = target.requests();
assert!(
requests
.iter()
.any(|record| record.operation == FakeTargetOperation::PutObject && record.proxy_headers.ssec_transport_present),
"the SSE-C probe PUT must carry the X-Rustfs-Replication-* transport headers; journal: {requests:?}"
);
// No probe residue, including the SSE-C probe version.
let probe_put = requests
.into_iter()
.find(|record| record.operation == FakeTargetOperation::PutObject)
.ok_or("the probe PUT never reached the fake target")?;
let probe_key = probe_put.key.ok_or("probe PUT journal record has no key")?;
assert!(
target.stored_versions(target_bucket, &probe_key).is_empty(),
"all probe versions must be cleaned up"
);
target.shutdown().await;
Ok(())
}
/// C1 (backlog#1675 P1-22): heal-path convergence for SSE-C. An SSE-C object
/// whose live replication failed during a target outage must converge through
/// the scanner/heal compensation once the target returns — passing the N2
/// HEAD-back audit against the recovered RustFS target — and the replica must
/// be readable with the customer key.
#[tokio::test]
#[serial]
async fn test_bucket_replication_sse_c_heals_after_target_outage() -> TestResult {
init_logging();
let (source_env, mut target_env, source_bucket, target_bucket) =
build_sse_replication_pair("ssec-heal", false, false).await?;
let source_client = source_env.create_s3_client();
let key = "ssec-heal-contract.txt";
let body = b"repl-22 ssec heal payload".to_vec();
let customer_key = BASE64_STANDARD.encode(REPL17_SSEC_KEY);
let customer_key_md5 = sse_customer_key_md5_base64(REPL17_SSEC_KEY);
// Target outage: the SSE-C write cannot replicate.
target_env.stop_server();
source_client
.put_object()
.bucket(&source_bucket)
.key(key)
.body(ByteStream::from(body.clone()))
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
// The failure is observable on the source (SSE-C HEAD needs the key).
let deadline = tokio::time::Instant::now() + Duration::from_secs(30);
loop {
let head = source_client
.head_object()
.bucket(&source_bucket)
.key(key)
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
match head.replication_status().map(|status| status.as_str()) {
Some("PENDING") | Some("FAILED") => break,
other => {
if tokio::time::Instant::now() >= deadline {
return Err(format!("source SSE-C object never reported PENDING/FAILED; last status={other:?}").into());
}
sleep(Duration::from_millis(200)).await;
}
}
}
// Recover the target in place; the source scanner re-drives the failure.
target_env
.restart_server_preserving_data(vec![], &[("NO_PROXY", "127.0.0.1,localhost"), ("HTTP_PROXY", ""), ("HTTPS_PROXY", "")])
.await?;
wait_for_source_replication_status(&source_client, &source_bucket, key, "COMPLETED", true).await?;
// The healed replica is a REPLICA (status surfaces on HEAD) readable with
// the customer key.
let target_client = target_env.create_s3_client();
let replica_head = target_client
.head_object()
.bucket(&target_bucket)
.key(key)
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
assert_eq!(
replica_head.replication_status().map(|status| status.as_str()),
Some("REPLICA"),
"the healed copy must carry REPLICA status"
);
let replica = target_client
.get_object()
.bucket(&target_bucket)
.key(key)
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
assert_eq!(replica.sse_customer_algorithm(), Some("AES256"));
assert_eq!(replica.body.collect().await?.into_bytes().as_ref(), body.as_slice());
Ok(())
}
/// C1 (backlog#1675 P1-22): existing-object resync for SSE-C. An SSE-C object
/// written BEFORE any replication config must reach the RustFS target through
/// the existing-object resync (`replicate_all` transport, N2-audited), land as
/// a REPLICA, and read back with the customer key.
#[tokio::test]
#[serial]
async fn test_bucket_replication_sse_c_existing_object_resync() -> TestResult {
init_logging();
let mut source_env = RustFSTestEnvironment::new().await?;
let mut source_process_env = replication_fast_env();
source_process_env.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
source_process_env.extend_from_slice(FAST_SCANNER_ENV);
source_process_env.extend_from_slice(&[("NO_PROXY", "127.0.0.1,localhost"), ("HTTP_PROXY", ""), ("HTTPS_PROXY", "")]);
source_env.start_rustfs_server_with_env(vec![], &source_process_env).await?;
let mut target_env = RustFSTestEnvironment::new().await?;
target_env
.start_rustfs_server_without_cleanup_with_env(&[
("NO_PROXY", "127.0.0.1,localhost"),
("HTTP_PROXY", ""),
("HTTPS_PROXY", ""),
])
.await?;
let source_bucket = "ssec-existing-src";
let target_bucket = "ssec-existing-dst";
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
target_client.create_bucket().bucket(target_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
enable_bucket_versioning(&target_env, target_bucket).await?;
// The SSE-C object exists before any replication wiring.
let key = "ssec-existing-contract.txt";
let body = b"repl-22 ssec existing-object payload".to_vec();
let customer_key = BASE64_STANDARD.encode(REPL17_SSEC_KEY);
let customer_key_md5 = sse_customer_key_md5_base64(REPL17_SSEC_KEY);
source_client
.put_object()
.bucket(source_bucket)
.key(key)
.body(ByteStream::from(body.clone()))
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
// Wire replication (existing-object enabled) and drive a resync.
let target_arn = set_replication_target(&source_env, source_bucket, &target_env, target_bucket).await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let (reset_arn, reset_id) = start_bucket_replication_reset(&source_env, source_bucket).await?;
assert_eq!(reset_arn, target_arn);
let terminal = wait_for_replication_reset_target(&source_env, source_bucket, &target_arn, |status| {
status.reset_id == reset_id && matches!(status.status.as_str(), "Completed" | "Failed")
})
.await?;
assert_eq!(terminal.status, "Completed", "SSE-C existing-object resync must complete");
assert!(terminal.replicated_count >= 1, "the existing SSE-C object must have been resynced");
// The replica is a REPLICA (status surfaces on HEAD) readable with the
// customer key.
let replica_head = target_client
.head_object()
.bucket(target_bucket)
.key(key)
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
assert_eq!(
replica_head.replication_status().map(|status| status.as_str()),
Some("REPLICA"),
"the resynced copy must carry REPLICA status"
);
let replica = target_client
.get_object()
.bucket(target_bucket)
.key(key)
.sse_customer_algorithm("AES256")
.sse_customer_key(&customer_key)
.sse_customer_key_md5(&customer_key_md5)
.send()
.await?;
assert_eq!(replica.sse_customer_algorithm(), Some("AES256"));
assert_eq!(replica.body.collect().await?.into_bytes().as_ref(), body.as_slice());
// No plaintext leak: the replica stays unreadable without the key.
assert!(
target_client
.get_object()
.bucket(target_bucket)
.key(key)
.send()
.await
.is_err(),
"SSE-C replica must not be readable without the customer key"
);
Ok(())
}
/// backlog#1147 repl-17 / backlog#1783: SSE-S3 objects replicate by decrypting
/// at the source and re-encrypting on the target with the target's own KMS.
/// The property backlog#1291 pinned — never a silent plaintext replica — still
@@ -8417,3 +8824,304 @@ async fn test_scanner_never_compensates_when_existing_object_replication_disable
Ok(())
}
/// Shared setup for the P1-5 read-proxy scenarios (backlog#1675): a RustFS
/// source with an enabled replication rule pointing at the fake target, and
/// an object seeded DIRECTLY on the target — it exists remotely but not
/// locally, exactly the active-active replication-lag window the read proxy
/// serves.
async fn start_read_proxy_lab(
source_bucket: &str,
target_bucket: &str,
) -> Result<(FakeS3Target, RustFSTestEnvironment, Client, Client), Box<dyn Error + Send + Sync>> {
let target = FakeS3Target::start().await?;
target.create_bucket(target_bucket);
target.assign_own_version_ids(true);
let mut source_env = RustFSTestEnvironment::new().await?;
let mut process_env = replication_fast_env();
process_env.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
process_env.extend_from_slice(&[
("NO_PROXY", "127.0.0.1,localhost"),
("HTTP_PROXY", ""),
("HTTPS_PROXY", ""),
("RUST_LOG", "error"),
]);
source_env.start_rustfs_server_with_env(vec![], &process_env).await?;
let source_client = source_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
let target_arn = set_replication_target_with_options(
&source_env,
source_bucket,
ReplicationTargetOptions {
endpoint: &target.address(),
access_key: FAKE_ACCESS_KEY,
secret_key: FAKE_SECRET_KEY,
target_bucket,
secure: false,
skip_tls_verify: false,
ca_cert_pem: None,
},
)
.await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let target_client = Client::from_conf(crate::common::build_test_s3_config(
target.endpoint(),
FAKE_ACCESS_KEY,
FAKE_SECRET_KEY,
None,
"read-proxy-e2e",
));
Ok((target, source_env, source_client, target_client))
}
/// P1-5 (backlog#1675): during the active-active replication lag window a
/// GET/HEAD for an object the local site does not have yet is proxied to the
/// replication target. Pins the wire contract: the anti-loop
/// `source-proxy-request` marker is sent, the replication worker's
/// `source-replication-check` SSE-C exemption is NEVER sent, client SSE-C
/// headers are forwarded verbatim, and an inbound request that was itself
/// proxied is answered locally (404) without touching the target.
#[tokio::test]
#[serial]
async fn test_get_and_head_proxy_unreplicated_object_to_replication_target() -> TestResult {
init_logging();
let source_bucket = "proxy-read-src";
let target_bucket = "proxy-read-dst";
let (target, source_env, source_client, target_client) = start_read_proxy_lab(source_bucket, target_bucket).await?;
let payload = b"proxy payload".to_vec();
target_client
.put_object()
.bucket(target_bucket)
.key("proxy-only")
.body(ByteStream::from(payload.clone()))
.send()
.await?;
target.take_requests();
// a. GET of the locally-missing object is served through the proxy.
let got = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.send()
.await
.map_err(|err| format!("proxied GET failed: {}", err.into_service_error()))?;
assert_eq!(got.content_length, Some(payload.len() as i64));
let body = got.body.collect().await?.into_bytes();
assert_eq!(body.as_ref(), payload.as_slice(), "proxied GET must stream the target's body");
let get_record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::GetObject && record.key.as_deref() == Some("proxy-only"))
.ok_or("fake target never received the proxied GET")?;
assert_eq!(
get_record.proxy_headers.source_proxy_request.as_deref(),
Some("true"),
"proxied GET must carry the anti-loop source-proxy-request marker"
);
assert!(
get_record.proxy_headers.replication_check.is_none(),
"proxied GET must never carry the replication worker's source-replication-check exemption"
);
assert!(
get_record.proxy_headers.ssec_algorithm.is_none() && !get_record.proxy_headers.ssec_key_present,
"no client SSE-C headers were sent, so none may be forwarded"
);
// a2. Client SSE-C headers travel verbatim to the target (the target owns
// the real SSE-C decryption; the plaintext fake simply ignores them).
target.take_requests();
let ssec_key = "01234567890123456789012345678901";
let ssec_key_b64 = BASE64_STANDARD.encode(ssec_key);
let ssec_key_md5 = sse_customer_key_md5_base64(ssec_key);
let _ = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.sse_customer_algorithm("AES256")
.sse_customer_key(&ssec_key_b64)
.sse_customer_key_md5(&ssec_key_md5)
.send()
.await
.map_err(|err| format!("proxied SSE-C GET failed: {}", err.into_service_error()))?;
let ssec_record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::GetObject && record.key.as_deref() == Some("proxy-only"))
.ok_or("fake target never received the proxied SSE-C GET")?;
assert_eq!(ssec_record.proxy_headers.ssec_algorithm.as_deref(), Some("AES256"));
assert!(ssec_record.proxy_headers.ssec_key_present, "SSE-C key header must be forwarded verbatim");
assert_eq!(ssec_record.proxy_headers.ssec_key_md5.as_deref(), Some(ssec_key_md5.as_str()));
assert!(ssec_record.proxy_headers.replication_check.is_none());
// b. HEAD of the locally-missing object is served through the proxy.
target.take_requests();
let head = source_client
.head_object()
.bucket(source_bucket)
.key("proxy-only")
.send()
.await
.map_err(|err| format!("proxied HEAD failed: {}", err.into_service_error()))?;
assert_eq!(head.content_length, Some(payload.len() as i64));
let head_record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::HeadObject && record.key.as_deref() == Some("proxy-only"))
.ok_or("fake target never received the proxied HEAD")?;
assert_eq!(head_record.proxy_headers.source_proxy_request.as_deref(), Some("true"));
assert!(head_record.proxy_headers.replication_check.is_none());
// c. Anti-loop: an inbound request that already carries the proxy marker
// is answered locally with 404 and never forwarded to the target.
target.take_requests();
let err = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.customize()
.mutate_request(|req| {
req.headers_mut().insert("x-minio-source-proxy-request", "true");
})
.send()
.await
.expect_err("anti-loop GET must fail locally instead of proxying");
let service_err = err.into_service_error();
assert!(service_err.is_no_such_key(), "anti-loop GET must 404, got: {service_err}");
assert!(
!target
.requests()
.iter()
.any(|record| record.operation == FakeTargetOperation::GetObject),
"anti-loop GET must not reach the replication target; journal: {:?}",
target.requests()
);
// c2. MinIO ProxyHeaderSet parity: the header's mere PRESENCE disables
// proxying — "false" is exactly what a peer's replication worker sends on
// its convergence HEADs, and proxying that miss back would fake
// convergence.
target.take_requests();
let err = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.customize()
.mutate_request(|req| {
req.headers_mut().insert("x-minio-source-proxy-request", "false");
})
.send()
.await
.expect_err("proxy-header-set GET must fail locally instead of proxying");
let service_err = err.into_service_error();
assert!(service_err.is_no_such_key(), "proxy-header-set GET must 404, got: {service_err}");
assert!(
!target
.requests()
.iter()
.any(|record| record.operation == FakeTargetOperation::GetObject),
"proxy-header-set GET must not reach the replication target; journal: {:?}",
target.requests()
);
// d. The replication worker's own convergence HEAD against the target
// must carry `source-proxy-request: false` (never proxied back) and the
// replication-check exemption. Trigger real replication and inspect the
// fake journal.
target.take_requests();
source_client
.put_object()
.bucket(source_bucket)
.key("worker-replicated")
.body(ByteStream::from_static(b"worker payload"))
.send()
.await?;
wait_for_target_request_version_id(&target, FakeTargetOperation::PutObject, "worker-replicated").await?;
let worker_head = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::HeadObject && record.key.as_deref() == Some("worker-replicated"))
.ok_or_else(|| format!("replication worker never HEAD-ed the target; journal: {:?}", target.requests()))?;
assert_eq!(
worker_head.proxy_headers.source_proxy_request.as_deref(),
Some("false"),
"worker convergence HEAD must send source-proxy-request: false so the target answers locally"
);
assert_eq!(
worker_head.proxy_headers.replication_check.as_deref(),
Some("true"),
"worker convergence HEAD keeps the replication-check exemption"
);
drop(source_env);
target.shutdown().await;
Ok(())
}
/// P1-5 (backlog#1675): GetObjectTagging for an object missing locally is
/// proxied to the replication target with the anti-loop marker, mirroring
/// MinIO `proxyGetTaggingToRepTarget`.
#[tokio::test]
#[serial]
async fn test_get_object_tagging_proxies_unreplicated_object_to_replication_target() -> TestResult {
init_logging();
let source_bucket = "proxy-tag-src";
let target_bucket = "proxy-tag-dst";
let (target, source_env, source_client, target_client) = start_read_proxy_lab(source_bucket, target_bucket).await?;
target_client
.put_object()
.bucket(target_bucket)
.key("proxy-tagged")
.body(ByteStream::from_static(b"tagged payload"))
.send()
.await?;
target_client
.put_object_tagging()
.bucket(target_bucket)
.key("proxy-tagged")
.tagging(
aws_sdk_s3::types::Tagging::builder()
.tag_set(aws_sdk_s3::types::Tag::builder().key("team").value("storage").build()?)
.build()?,
)
.send()
.await?;
target.take_requests();
let tags = source_client
.get_object_tagging()
.bucket(source_bucket)
.key("proxy-tagged")
.send()
.await
.map_err(|err| format!("proxied GetObjectTagging failed: {}", err.into_service_error()))?;
assert_eq!(tags.tag_set.len(), 1, "proxied tagging read must return the target's tags");
assert_eq!(tags.tag_set[0].key.as_str(), "team");
assert_eq!(tags.tag_set[0].value.as_str(), "storage");
let record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::GetObjectTagging && record.key.as_deref() == Some("proxy-tagged"))
.ok_or("fake target never received the proxied GetObjectTagging")?;
assert_eq!(
record.proxy_headers.source_proxy_request.as_deref(),
Some("true"),
"proxied tagging read must carry the anti-loop marker"
);
assert!(record.proxy_headers.replication_check.is_none());
drop(source_env);
target.shutdown().await;
Ok(())
}
+8 -7
View File
@@ -32,7 +32,7 @@ pub mod bucket {
pub mod bucket_target_sys {
pub use crate::bucket::bucket_target_sys::{
AdvancedPutOptions, BucketTargetError, BucketTargetSys, PutObjectOptions, RemoveObjectOptions, S3ClientError,
TargetClient, append_version_id_query,
SsecPassthroughCapability, TargetClient, append_version_id_query,
};
}
@@ -198,12 +198,13 @@ pub mod bucket {
ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog, TargetReplicationResyncStatus,
VersionPurgeStatusType, XferStats, commit_force_delete_intent, complete_force_delete_intent,
delete_replication_state_from_config, delete_replication_version_id, get_global_replication_pool,
get_global_replication_stats, init_background_replication, invalid_replication_config_status_field,
persist_force_delete_intent, read_durable_mrf_backlog, replication_state_to_filemeta, replication_status_to_filemeta,
replication_statuses_map, replication_target_arns, resync_start_conflict_id, should_remove_replication_target,
should_schedule_delete_replication, should_use_existing_delete_replication_info,
should_use_existing_delete_replication_source, unsupported_replication_config_field,
validate_replication_config_structure, validate_replication_config_target_arns, version_purge_status_to_filemeta,
get_global_replication_stats, get_proxy_targets, init_background_replication,
invalid_replication_config_status_field, persist_force_delete_intent, read_durable_mrf_backlog,
replication_state_to_filemeta, replication_status_to_filemeta, replication_statuses_map, replication_target_arns,
resync_start_conflict_id, should_remove_replication_target, should_schedule_delete_replication,
should_use_existing_delete_replication_info, should_use_existing_delete_replication_source,
unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns,
version_purge_status_to_filemeta,
};
}
+316 -4
View File
@@ -27,10 +27,15 @@ use aws_sdk_s3::config::SharedHttpClient;
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::error::SdkError;
use aws_sdk_s3::operation::complete_multipart_upload::CompleteMultipartUploadOutput;
use aws_sdk_s3::operation::delete_object_tagging::{DeleteObjectTaggingError, DeleteObjectTaggingOutput};
use aws_sdk_s3::operation::get_object::{GetObjectError, GetObjectOutput};
use aws_sdk_s3::operation::get_object_tagging::{GetObjectTaggingError, GetObjectTaggingOutput};
use aws_sdk_s3::operation::head_bucket::HeadBucketError;
use aws_sdk_s3::operation::head_object::HeadObjectError;
use aws_sdk_s3::operation::put_object_tagging::{PutObjectTaggingError, PutObjectTaggingOutput};
use aws_sdk_s3::operation::upload_part::UploadPartOutput;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::Tagging as SdkTagging;
use aws_sdk_s3::types::{
ChecksumMode, CompletedMultipartUpload, CompletedPart, ObjectLockLegalHoldStatus, ObjectLockRetentionMode,
};
@@ -57,8 +62,8 @@ use rustfs_utils::http::{
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_REPLICATION_CHECK,
SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_PROXY_REQUEST,
SUFFIX_SOURCE_REPLICATION_CHECK, 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,
};
@@ -294,9 +299,51 @@ struct TargetClientBuildProbe {
release: Arc<tokio::sync::Semaphore>,
}
/// Whether a replication target preserves the SSE-C passthrough transport
/// headers (`X-Rustfs-Replication-*`) end to end.
///
/// A target that silently drops those headers (MinIO, generic S3) stores the
/// forwarded ciphertext without its decryption material — an unreadable
/// replica that used to report COMPLETED. The replication worker audits the
/// first passthrough PUT per target (HEAD-back for SSE-C evidence) and caches
/// the verdict here; a fresh `Unsupported` fails SSE-C replication closed
/// before any PUT is sent. Entries follow the `arn_remotes_map` lifecycle
/// (rebuilding or removing a target resets its capability to `Unknown`) and
/// additionally expire after [`SSEC_PASSTHROUGH_CAPABILITY_TTL`], after which
/// the next attempt re-audits.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SsecPassthroughCapability {
#[default]
Unknown,
Supported,
Unsupported,
}
/// How long an audited SSE-C passthrough verdict stays authoritative.
///
/// Trade-off: without a TTL a verdict is sticky for the process lifetime —
/// an `Unsupported` target that gets upgraded (or re-probed only via
/// replication-check) would keep failing SSE-C replication forever, and the
/// fail-open twin: a `Supported` verdict would outlive a backend swapped
/// behind the same endpoint/ARN. With the TTL, a bad target costs at most
/// one wasted PUT+HEAD audit per TTL window, and a changed backend is
/// re-discovered within the same window.
pub const SSEC_PASSTHROUGH_CAPABILITY_TTL: Duration = Duration::from_secs(10 * 60);
/// A recorded SSE-C passthrough verdict plus when it was recorded, so reads
/// can report staleness against [`SSEC_PASSTHROUGH_CAPABILITY_TTL`].
#[derive(Debug, Clone, Copy)]
struct SsecPassthroughRecord {
capability: SsecPassthroughCapability,
recorded_at: Instant,
}
#[derive(Debug, Default)]
pub struct BucketTargetSys {
pub arn_remotes_map: Arc<RwLock<HashMap<String, ArnTarget>>>,
/// SSE-C passthrough capability verdicts keyed by target ARN. See
/// [`SsecPassthroughCapability`]; reset alongside `arn_remotes_map`.
ssec_passthrough_map: Arc<RwLock<HashMap<String, SsecPassthroughRecord>>>,
pub targets_map: Arc<RwLock<HashMap<String, Vec<BucketTarget>>>>,
pub h_mutex: Arc<RwLock<HashMap<String, EpHealth>>>,
target_h_mutex: Arc<RwLock<HashMap<String, EpHealth>>>,
@@ -317,6 +364,7 @@ impl BucketTargetSys {
fn new() -> Self {
Self {
arn_remotes_map: Arc::new(RwLock::new(HashMap::new())),
ssec_passthrough_map: Arc::new(RwLock::new(HashMap::new())),
targets_map: Arc::new(RwLock::new(HashMap::new())),
h_mutex: Arc::new(RwLock::new(HashMap::new())),
target_h_mutex: Arc::new(RwLock::new(HashMap::new())),
@@ -580,19 +628,59 @@ impl BucketTargetSys {
let update_mutex = self.target_update_mutex(bucket).await;
let _update_guard = update_mutex.lock().await;
// Lock order: targets_map, then arn_remotes_map, then target_h_mutex.
// Lock order: targets_map, then arn_remotes_map, then target_h_mutex,
// then ssec_passthrough_map (always last; also taken standalone by the
// capability accessors).
let mut targets_map = self.targets_map.write().await;
let mut arn_remotes_map = self.arn_remotes_map.write().await;
let mut health_map = self.target_h_mutex.write().await;
if let Some(targets) = targets_map.remove(bucket) {
let mut ssec_map = self.ssec_passthrough_map.write().await;
for target in targets {
arn_remotes_map.remove(&target.arn);
health_map.remove(&target.arn);
ssec_map.remove(&target.arn);
}
}
}
/// Cached SSE-C passthrough capability for a target ARN, plus whether the
/// verdict is older than [`SSEC_PASSTHROUGH_CAPABILITY_TTL`]. `(Unknown,
/// false)` when no verdict has been recorded since the target was built.
/// Staleness is computed here so the gate policy stays a pure function.
pub async fn ssec_passthrough_capability(&self, arn: &str) -> (SsecPassthroughCapability, bool) {
match self.ssec_passthrough_map.read().await.get(arn) {
Some(record) => (record.capability, record.recorded_at.elapsed() >= SSEC_PASSTHROUGH_CAPABILITY_TTL),
None => (SsecPassthroughCapability::Unknown, false),
}
}
/// Record an audited SSE-C passthrough verdict for a target ARN. Written by
/// the replication worker's HEAD-back audit and by the replication-check
/// SsecPassthrough probe phase.
pub async fn record_ssec_passthrough_capability(&self, arn: &str, capability: SsecPassthroughCapability) {
self.ssec_passthrough_map.write().await.insert(
arn.to_string(),
SsecPassthroughRecord {
capability,
recorded_at: Instant::now(),
},
);
}
/// Test hook: age an existing verdict so TTL expiry is observable without
/// waiting out the real window.
#[cfg(test)]
pub(crate) async fn backdate_ssec_passthrough_capability(&self, arn: &str, age: Duration) {
let backdated = Instant::now()
.checked_sub(age)
.expect("system uptime must exceed the backdate age");
if let Some(record) = self.ssec_passthrough_map.write().await.get_mut(arn) {
record.recorded_at = backdated;
}
}
pub async fn set_target(
&self,
bucket: &str,
@@ -948,15 +1036,21 @@ impl BucketTargetSys {
}
}
// Lock order: targets_map, then arn_remotes_map, then target_h_mutex.
// Lock order: targets_map, then arn_remotes_map, then target_h_mutex,
// then ssec_passthrough_map (always last; also taken standalone by the
// capability accessors).
let mut targets_map = self.targets_map.write().await;
let mut arn_remotes_map = self.arn_remotes_map.write().await;
let mut health_map = self.target_h_mutex.write().await;
// Remove existing targets
if let Some(existing_targets) = targets_map.remove(bucket) {
let mut ssec_map = self.ssec_passthrough_map.write().await;
for target in existing_targets {
arn_remotes_map.remove(&target.arn);
health_map.remove(&target.arn);
// A rebuilt/edited target may point at a different service:
// the SSE-C passthrough verdict must be re-audited from Unknown.
ssec_map.remove(&target.arn);
self.update_bandwidth_limit(bucket, &target.arn, 0);
}
}
@@ -1446,6 +1540,43 @@ fn resolve_put_api_version_id(source_version_id: &str) -> Option<&str> {
}
}
/// Resolve the S3 `versionId` for a proxied read against a remote target.
/// RustFS represents the null version internally as the nil UUID while the S3
/// API addresses it as the literal "null" (same mapping as
/// [`resolve_put_api_version_id`]); empty means "no version requested".
pub(crate) fn resolve_read_api_version_id(version_id: Option<String>) -> Option<String> {
let version_id = version_id?;
let trimmed = version_id.trim();
if trimmed.is_empty() {
None
} else if Uuid::parse_str(trimmed).is_ok_and(|uuid| uuid.is_nil()) {
Some(rustfs_filemeta::NULL_VERSION_ID.to_string())
} else {
Some(trimmed.to_string())
}
}
/// Outbound header set for a proxied read: the caller-provided passthrough
/// headers (client SSE-C key family, conditional headers) plus the anti-loop
/// `source-proxy-request` marker in both the x-rustfs- and x-minio- prefixes
/// (a MinIO target only understands the latter). Never adds
/// `source-replication-check`: that exemption channel belongs exclusively to
/// the replication worker's HEAD.
fn proxy_outbound_headers(mut extra_headers: HeaderMap) -> HeaderMap {
insert_header(&mut extra_headers, SUFFIX_SOURCE_PROXY_REQUEST, "true");
extra_headers
}
/// Copy `headers` onto an SDK request inside `customize().map_request` (runs
/// before signing, so the headers join the SigV4 canonical request).
fn apply_extra_headers(mut req: HttpRequest, headers: &HeaderMap) -> Result<HttpRequest, std::convert::Infallible> {
for (k, v) in headers.iter() {
req.headers_mut()
.insert(k.as_str().to_string(), v.to_str().unwrap_or("").to_string());
}
Ok(req)
}
/// Append `versionId=<id>` to an already-built request URI. aws-sdk-s3's
/// `PutObjectInput` / `CreateMultipartUploadInput` expose no version id
/// member, so the query is spliced in via `map_request`, which runs at
@@ -1851,6 +1982,13 @@ impl TargetClient {
// worker cannot hold; otherwise SSE-C replicas never converge on HEAD.
let mut headers = HeaderMap::new();
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_CHECK, "true");
// `source-proxy-request: false` (MinIO `ProxyHeaderSet` semantics):
// the header's mere presence tells the receiver to answer LOCALLY
// instead of proxying the miss back to us. Without it, a not-found on
// the target gets read-proxied back to this source, echoes the source
// object with an identical ETag, and the worker concludes the object
// already converged — so it never actually replicates it.
insert_header(&mut headers, SUFFIX_SOURCE_PROXY_REQUEST, "false");
match self
.client
.head_object()
@@ -1875,6 +2013,129 @@ impl TargetClient {
}
}
/// HEAD used by the read-proxy path (GET/HEAD of an object not yet
/// replicated locally, MinIO `proxyHeadToRepTarget`).
///
/// Deliberately different from [`TargetClient::head_object`]: it must NOT
/// send `source-replication-check` — that header is the replication
/// worker's SSE-C metadata exemption channel. A proxied client request
/// instead forwards the client's own SSE-C headers (`extra_headers`) so
/// the target performs the real SSE-C validation/decryption. The
/// `source-proxy-request` marker is always added so the target does not
/// proxy the request onward (anti-loop).
pub async fn head_object_for_proxy(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
range: Option<String>,
part_number: Option<i32>,
extra_headers: HeaderMap,
) -> Result<HeadObjectOutput, SdkError<HeadObjectError>> {
let headers = proxy_outbound_headers(extra_headers);
self.client
.head_object()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.set_range(range)
.set_part_number(part_number)
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// GET used by the read-proxy path (MinIO `proxyGetToReplicationTarget`).
/// Returns the streaming SDK output; callers must forward the body without
/// buffering it. Same header contract as [`Self::head_object_for_proxy`]:
/// anti-loop marker on, replication-check never sent, client SSE-C /
/// conditional headers forwarded verbatim via `extra_headers`.
pub async fn get_object(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
range: Option<String>,
part_number: Option<i32>,
extra_headers: HeaderMap,
) -> Result<GetObjectOutput, SdkError<GetObjectError>> {
let headers = proxy_outbound_headers(extra_headers);
self.client
.get_object()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.set_range(range)
.set_part_number(part_number)
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// GetObjectTagging for the tagging read-proxy path
/// (MinIO `proxyGetTaggingToRepTarget`). Anti-loop marker always added.
pub async fn get_object_tagging(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
) -> Result<GetObjectTaggingOutput, SdkError<GetObjectTaggingError>> {
let headers = proxy_outbound_headers(HeaderMap::new());
self.client
.get_object_tagging()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// PutObjectTagging for the tagging proxy path
/// (MinIO `proxyTaggingToRepTarget`). Anti-loop marker always added.
pub async fn put_object_tagging(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
tagging: SdkTagging,
) -> Result<PutObjectTaggingOutput, SdkError<PutObjectTaggingError>> {
let headers = proxy_outbound_headers(HeaderMap::new());
self.client
.put_object_tagging()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.tagging(tagging)
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// DeleteObjectTagging for the tagging proxy path
/// (MinIO `proxyTaggingToRepTarget`). Anti-loop marker always added.
pub async fn delete_object_tagging(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
) -> Result<DeleteObjectTaggingOutput, SdkError<DeleteObjectTaggingError>> {
let headers = proxy_outbound_headers(HeaderMap::new());
self.client
.delete_object_tagging()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// On success returns the version id the target assigned (from
/// `x-amz-version-id`), letting callers audit the version-identity
/// contract — a target that adopts the source version echoes it back.
@@ -2504,6 +2765,57 @@ mod tests {
assert_eq!(health.last_online, Some(now));
}
/// N2 TTL contract, both flip directions: a recorded verdict is fresh
/// until [`SSEC_PASSTHROUGH_CAPABILITY_TTL`], then reads as expired; a
/// re-audit that records the OPPOSITE verdict replaces it as fresh. The
/// worker gate maps expired verdicts to ProceedWithAudit (pinned in
/// `replication_target_boundary`), so together this proves an Unsupported
/// target recovers to Supported through the audit once its verdict ages
/// out — and a stale Supported one is re-proven rather than trusted.
#[tokio::test]
async fn ssec_passthrough_capability_ttl_expires_and_reaudit_flips_verdict() {
let sys = BucketTargetSys::default();
let arn = "arn:rustfs:replication:us-east-1:bucket:ssec-ttl";
let expired_age = SSEC_PASSTHROUGH_CAPABILITY_TTL + Duration::from_secs(1);
assert_eq!(
sys.ssec_passthrough_capability(arn).await,
(SsecPassthroughCapability::Unknown, false),
"an unrecorded target must read Unknown and never expired"
);
sys.record_ssec_passthrough_capability(arn, SsecPassthroughCapability::Unsupported)
.await;
assert_eq!(
sys.ssec_passthrough_capability(arn).await,
(SsecPassthroughCapability::Unsupported, false)
);
sys.backdate_ssec_passthrough_capability(arn, expired_age).await;
assert_eq!(
sys.ssec_passthrough_capability(arn).await,
(SsecPassthroughCapability::Unsupported, true),
"an aged-out Unsupported verdict must read expired so the gate re-audits"
);
// The re-audit against an upgraded target records Supported afresh.
sys.record_ssec_passthrough_capability(arn, SsecPassthroughCapability::Supported)
.await;
assert_eq!(
sys.ssec_passthrough_capability(arn).await,
(SsecPassthroughCapability::Supported, false),
"a fresh Supported verdict replaces the expired Unsupported one"
);
// And the fail-open twin: Supported also ages out.
sys.backdate_ssec_passthrough_capability(arn, expired_age).await;
assert_eq!(
sys.ssec_passthrough_capability(arn).await,
(SsecPassthroughCapability::Supported, true),
"an aged-out Supported verdict must read expired so the gate re-proves it"
);
}
#[tokio::test]
async fn list_targets_applies_health_stats_by_arn_and_preserves_endpoint_port() {
let sys = BucketTargetSys::default();
@@ -14,6 +14,7 @@ paths.
| `datatypes.rs` | ECStore compatibility re-export for resync status enums. | Re-exports `rustfs-replication` contracts while downstream facade consumers migrate. |
| `replication_object_decision_boundary.rs` | Object replication option DTOs, resync target projection, delete replication decisions, and multipart planning helpers. | Keeps ECStore runtime modules from importing object decision contracts directly from `rustfs-replication`. |
| `replication_pool.rs` | Replication queue, worker pool, MRF persistence, bucket stats, and delete/object scheduling. | Depends on bucket target sys, bucket metadata sys, metadata paths, queue contracts through the queue boundary, file metadata replication contracts through local boundaries, config storage, storage contracts through the replication storage boundary, runtime sources, and notification state. |
| `replication_proxy.rs` | Proxy-target selection for GET/HEAD/Tagging reads of objects not yet replicated locally (MinIO `getProxyTargets` parity: anti-loop, version-suspended, and no-config empty branches). | Uses replication config lookup, rule matching, and target clients through local boundaries. |
| `replication_queue_boundary.rs` | Queue/admission DTOs, heal queue DTOs, worker sizing, and backpressure helpers. | Keeps ECStore runtime modules from importing queue/backpressure contracts directly from `rustfs-replication`. |
| `replication_resync_boundary.rs` | Resync DTOs, status classifiers, persisted resync/MRF codec wrappers, and ECStore error mapping. | Keeps ECStore runtime modules from importing resync contract helpers directly from `rustfs-replication`. |
| `replication_resyncer.rs` | Object replication, delete replication, resync execution, target calls, and multipart target upload paths. | Depends on target calls and target config types through the replication target boundary, metadata paths and metadata systems through the replication metadata boundary, file metadata replication contracts through the filemeta boundary, object decisions and multipart planning through the object decision boundary, resync contracts through the resync boundary, queue DTOs through the queue boundary, error contracts through the error boundary, versioning systems, storage contracts through the replication storage boundary, config-derived storage class labels through the config store, runtime sources, notification events and local event host selection through the event sink, bandwidth reader wrapping, and SetDisks lock timing. |
@@ -29,6 +29,7 @@ mod replication_object_bridge;
mod replication_object_config;
mod replication_object_decision_boundary;
pub(crate) mod replication_pool;
mod replication_proxy;
mod replication_queue_boundary;
mod replication_resync_boundary;
mod replication_resyncer;
@@ -74,6 +75,7 @@ pub use replication_pool::{
get_global_replication_pool, get_global_replication_stats, init_background_replication, persist_force_delete_intent,
read_durable_mrf_backlog, resync_start_conflict_id,
};
pub use replication_proxy::get_proxy_targets;
pub use replication_queue_boundary::{
DeletedObjectReplicationInfo, ReplicationBatchAdmission, ReplicationHealQueueResult, ReplicationOperation,
ReplicationPriority, ReplicationQueueAdmission,
@@ -0,0 +1,150 @@
// 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.
//! Proxy-target selection for reads of objects not yet replicated locally
//! (MinIO `getProxyTargets`, bucket-replication.go).
//!
//! During the active-active replication lag window a GET/HEAD/Tagging request
//! for an object the local site does not have yet may be served by proxying to
//! a replication target. This module only *selects* the candidate targets; the
//! request-path callers perform the remote calls and response translation.
use std::sync::Arc;
use tracing::debug;
use super::replication_config_boundary::{ObjectOpts, ReplicationConfigurationExt as _};
use super::replication_object_config::get_replication_config;
use super::replication_storage_boundary::ObjectOptions;
use super::replication_target_boundary::{ReplicationTargetStore, TargetClient};
/// Returns the replication-target clients eligible to serve a proxied read of
/// `bucket/object`, in rule order. Mirrors MinIO's `getProxyTargets`:
///
/// - the `source-proxy-request` header family was present at all
/// (`opts.proxy_request` / `opts.proxy_header_set`, MinIO `ProxyRequest` /
/// `ProxyHeaderSet`) -> empty. "true" is the anti-loop marker of an
/// already-proxied client read; "false" is what a peer's replication
/// worker sends on convergence HEADs so the receiver answers locally —
/// proxying that miss back would echo the source object and fake
/// convergence, permanently skipping replication;
/// - the bucket's versioning is suspended for the object -> empty;
/// - no replication configuration / no matching rule -> empty;
/// - otherwise every distinct target ARN whose rules match the object,
/// resolved through the bucket target system, skipping targets that opted
/// out of proxying (`disable_proxy`).
pub async fn get_proxy_targets(bucket: &str, object: &str, opts: &ObjectOptions) -> Vec<Arc<TargetClient>> {
if opts.proxy_request || opts.proxy_header_set {
return Vec::new();
}
if opts.version_suspended {
return Vec::new();
}
let cfg = match get_replication_config(bucket).await {
Ok(Some(cfg)) => cfg,
Ok(None) => return Vec::new(),
Err(err) => {
debug!(bucket, object, error = %err, "read proxy: failed to load replication config; not proxying");
return Vec::new();
}
};
let arns = cfg.filter_target_arns(&ObjectOpts {
name: object.to_string(),
..Default::default()
});
let mut targets = Vec::with_capacity(arns.len());
for arn in arns {
let Some(client) = ReplicationTargetStore::remote_target_client(bucket, &arn).await else {
debug!(bucket, object, arn, "read proxy: no client for replication target ARN");
continue;
};
if client.disable_proxy {
continue;
}
targets.push(client);
}
targets
}
#[cfg(test)]
mod tests {
use super::*;
fn opts() -> ObjectOptions {
ObjectOptions::default()
}
/// Anti-loop: a request that was already proxied by a peer must never be
/// proxied onward, regardless of replication configuration.
#[tokio::test]
async fn proxy_request_yields_no_targets() {
let targets = get_proxy_targets(
"bucket",
"object",
&ObjectOptions {
proxy_request: true,
..opts()
},
)
.await;
assert!(targets.is_empty());
}
/// MinIO `ProxyHeaderSet` parity: the header family being present at all
/// disables proxying, even with the value "false" — that is what a
/// peer's replication worker sends on convergence HEADs.
#[tokio::test]
async fn proxy_header_set_yields_no_targets() {
let targets = get_proxy_targets(
"bucket",
"object",
&ObjectOptions {
proxy_header_set: true,
proxy_request: false,
..opts()
},
)
.await;
assert!(targets.is_empty());
}
/// Suspended versioning disables proxying (MinIO parity): the local null
/// version is authoritative and a remote read could resurrect data.
#[tokio::test]
async fn version_suspended_yields_no_targets() {
let targets = get_proxy_targets(
"bucket",
"object",
&ObjectOptions {
version_suspended: true,
..opts()
},
)
.await;
assert!(targets.is_empty());
}
/// A bucket without replication configuration has nothing to proxy to.
/// (No metadata system is running in unit tests, so the config lookup
/// resolves to "no configuration" — the same empty-result contract.)
#[tokio::test]
async fn missing_replication_config_yields_no_targets() {
let targets = get_proxy_targets("bucket-without-replication", "object", &opts()).await;
assert!(targets.is_empty());
}
}
@@ -15,10 +15,14 @@
use super::replication_error_boundary::{Error, Result};
use super::replication_filemeta_boundary::MrfReplicateEntry;
/// Kept test-only: the runtime consumer was the worker HEAD's fake proxy
/// counting (removed in backlog#1675 P1-5); the resyncer tests still pin the
/// classifier's semantics for the real client read-proxy failure accounting.
#[cfg(test)]
pub(crate) use rustfs_replication::should_count_head_proxy_failure;
pub use rustfs_replication::{BucketReplicationResyncStatus, ResyncOpts, ResyncStatusType, TargetReplicationResyncStatus};
pub(crate) use rustfs_replication::{
is_version_id_mismatch, resync_state_accepts_update, sanitize_resync_error_detail, should_auto_resume_resync,
should_count_head_proxy_failure,
};
#[allow(
@@ -36,6 +36,7 @@ use super::replication_object_decision_boundary::{
};
use super::replication_queue_boundary::{DeletedObjectReplicationInfo, ReplicationQueueAdmission};
use super::replication_resync_boundary::ResyncStatusType;
#[cfg(test)]
use super::replication_resync_boundary::should_count_head_proxy_failure;
use super::replication_resync_boundary::{
BucketReplicationResyncStatus, ResyncOpts, TargetReplicationResyncStatus, encode_resync_file, is_version_id_mismatch,
@@ -48,10 +49,12 @@ use super::replication_storage_boundary::{
ReplicationDeletedObject, ReplicationObjectIO, ReplicationStorage, StatObjectOptions, StorageObjectInfoOrErr, WalkOptions,
};
use super::replication_target_boundary::{
PutObjectOptions, PutObjectPartOptions, ReplicationTargetStore, TargetClient, replication_action_for_target_head,
ERR_REPLICATION_SSEC_PASSTHROUGH_UNSUPPORTED, PutObjectOptions, PutObjectPartOptions, ReplicationTargetStore,
SsecPassthroughCapability, SsecPassthroughGate, TargetClient, replication_action_for_target_head,
replication_complete_multipart_options, replication_delete_marker_purge_remove_options, replication_delete_remove_options,
replication_force_delete_remove_options, replication_object_is_ssec_encrypted, replication_put_object_header_size,
replication_put_object_options, replication_target_head_is_newer_null_version,
replication_put_object_options, replication_target_head_is_newer_null_version, resolve_read_api_version_id,
ssec_passthrough_evidence_present, ssec_passthrough_gate,
};
use super::replication_versioning_boundary::ReplicationVersioningStore;
use super::runtime_boundary as runtime_sources;
@@ -234,32 +237,18 @@ fn audit_target_version_identity(tgt_client: &TargetClient, source_version_id: &
}
}
fn is_head_proxy_failure(err: &SdkError<HeadObjectError>) -> bool {
let (is_not_found, code) = err
.as_service_error()
.map(|service_err| (service_err.is_not_found(), service_err.code()))
.unwrap_or((false, None));
let raw_status = err.raw_response().map(|resp| resp.status().as_u16());
should_count_head_proxy_failure(is_not_found, code, raw_status)
}
async fn record_proxy_request(bucket: &str, api: &str, is_err: bool) {
if let Some(stats) = runtime_sources::replication_stats() {
stats.inc_proxy(bucket, api, is_err).await;
}
}
async fn head_object_with_proxy_stats(
source_bucket: &str,
/// HEAD against a replication target on behalf of the replication worker
/// (resync/heal/delete convergence checks). This is NOT a client read proxy:
/// it must not touch the proxy metrics — those count only real GET/HEAD/
/// Tagging requests proxied for clients (see `replication_proxy.rs` /
/// `TargetClient::head_object_for_proxy`).
async fn head_object_for_worker(
target_client: &TargetClient,
target_bucket: &str,
object: &str,
version_id: Option<String>,
) -> std::result::Result<HeadObjectOutput, SdkError<HeadObjectError>> {
let result = target_client.head_object(target_bucket, object, version_id).await;
let is_err = result.as_ref().err().is_some_and(is_head_proxy_failure);
record_proxy_request(source_bucket, "HeadObject", is_err).await;
result
target_client.head_object(target_bucket, object, version_id).await
}
fn is_version_id_format_mismatch(err: &SdkError<HeadObjectError>) -> bool {
@@ -282,17 +271,124 @@ async fn mark_replication_target_offline_if_needed(target_client: &Arc<TargetCli
}
async fn head_object_fallback(
source_bucket: &str,
tgt_client: &TargetClient,
object: &str,
) -> std::result::Result<Option<HeadObjectOutput>, SdkError<HeadObjectError>> {
match head_object_with_proxy_stats(source_bucket, tgt_client, &tgt_client.bucket, object, None).await {
match head_object_for_worker(tgt_client, &tgt_client.bucket, object, None).await {
Ok(oi) => Ok(Some(oi)),
Err(e) if e.as_service_error().is_some_and(|se| se.is_not_found()) || has_raw_status(&e, 404) => Ok(None),
Err(e) => Err(e),
}
}
/// Resolve the N2 fail-closed gate for an SSE-C passthrough attempt against
/// this target. Returns `Some(audit_required)` when replication may proceed;
/// on a freshly-flagged header-dropping target it settles `rinfo` as FAILED
/// (no PUT is ever sent — the object stays on the normal MRF retry channel
/// and re-audits once the verdict's TTL expires or replication-check
/// re-probes the target) and returns `None`.
async fn resolve_ssec_passthrough_gate(
ssec: bool,
tgt_client: &TargetClient,
bucket: &str,
object: &str,
rinfo: &mut ReplicatedTargetInfo,
) -> Option<bool> {
let (capability, expired) = ReplicationTargetStore::ssec_passthrough_capability(&tgt_client.arn).await;
match ssec_passthrough_gate(ssec, capability, expired) {
SsecPassthroughGate::Proceed => Some(false),
SsecPassthroughGate::ProceedWithAudit => Some(true),
SsecPassthroughGate::FailClosed => {
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(ERR_REPLICATION_SSEC_PASSTHROUGH_UNSUPPORTED.to_string());
warn!(
event = EVENT_RESYNC_TARGET_OPERATION_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
bucket = %bucket,
object = %object,
arn = %tgt_client.arn,
operation = "ssec_passthrough_gate",
error = ERR_REPLICATION_SSEC_PASSTHROUGH_UNSUPPORTED,
"Replication target operation failed"
);
None
}
}
}
/// Judge SSE-C passthrough evidence on a HEAD of the replica and record the
/// capability verdict for the target. Returns true when the SSE-C material
/// provably survived; otherwise records `Unsupported` and settles `rinfo` as
/// FAILED so the attempt never reports a silently unreadable COMPLETED.
async fn settle_ssec_passthrough_evidence(
head: &HeadObjectOutput,
tgt_client: &TargetClient,
bucket: &str,
object: &str,
rinfo: &mut ReplicatedTargetInfo,
) -> bool {
if ssec_passthrough_evidence_present(head) {
ReplicationTargetStore::record_ssec_passthrough_capability(&tgt_client.arn, SsecPassthroughCapability::Supported).await;
return true;
}
ReplicationTargetStore::record_ssec_passthrough_capability(&tgt_client.arn, SsecPassthroughCapability::Unsupported).await;
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(ERR_REPLICATION_SSEC_PASSTHROUGH_UNSUPPORTED.to_string());
warn!(
event = EVENT_RESYNC_TARGET_OPERATION_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
bucket = %bucket,
object = %object,
arn = %tgt_client.arn,
endpoint = %tgt_client.endpoint,
operation = "ssec_passthrough_audit",
error = ERR_REPLICATION_SSEC_PASSTHROUGH_UNSUPPORTED,
"Replication target operation failed"
);
false
}
/// Post-PUT HEAD-back audit for an SSE-C passthrough replica, over the worker
/// HEAD channel (replication-check exemption plus the `source-proxy-request:
/// false` suppression header, so the target answers locally without a
/// customer key). A HEAD transport failure leaves the capability `Unknown`
/// but still fails this attempt: an unverifiable SSE-C replica must not
/// report COMPLETED.
async fn audit_ssec_passthrough_replica(
tgt_client: &Arc<TargetClient>,
bucket: &str,
object: &str,
version_id: Option<String>,
rinfo: &mut ReplicatedTargetInfo,
) -> bool {
// Address the replica the way the PUT named it: a nil source version id
// (versioning-suspended / null-version objects) maps to the "null"
// version, so the audit HEAD does not 4xx-loop on those objects.
let version_id = resolve_read_api_version_id(version_id);
match head_object_for_worker(tgt_client.as_ref(), &tgt_client.bucket, object, version_id).await {
Ok(head) => settle_ssec_passthrough_evidence(&head, tgt_client, bucket, object, rinfo).await,
Err(e) => {
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(format!("SSE-C passthrough audit HEAD failed: {e}"));
warn!(
event = EVENT_RESYNC_TARGET_OPERATION_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
bucket = %bucket,
object = %object,
arn = %tgt_client.arn,
operation = "ssec_passthrough_audit_head",
error = %e,
"Replication target operation failed"
);
mark_replication_target_offline_if_needed(tgt_client, &e).await;
false
}
}
}
static RESYNC_WORKER_COUNT: usize = 10;
fn resync_status_duration(
@@ -1186,7 +1282,7 @@ async fn verify_resync_head_result(
// (400). Re-verify without the versionId before
// concluding the object failed to replicate, instead
// of counting a well-replicated object as failed.
match head_object_fallback(&roi.bucket, target_client.as_ref(), &roi.name).await {
match head_object_fallback(target_client.as_ref(), &roi.name).await {
Ok(Some(_)) => {
st.replicated_count += 1;
st.replicated_size += roi.size;
@@ -1236,8 +1332,7 @@ async fn resync_worker_process_object<S: ReplicationStorage>(
let reset_id = target_client.reset_id.clone();
let head_result = head_object_with_proxy_stats(
bucket_name,
let head_result = head_object_for_worker(
target_client.as_ref(),
&target_client.bucket,
&roi.name,
@@ -2521,8 +2616,7 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli
let version_id = target_delete_version_id(version_id, is_version_purge);
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
match head_object_with_proxy_stats(
&dobj.bucket,
match head_object_for_worker(
tgt_client.as_ref(),
&tgt_client.bucket,
&dobj.delete_object.object_name,
@@ -2888,6 +2982,21 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
return rinfo;
}
// N2 fail-closed: never PUT SSE-C ciphertext at a target known to drop
// the passthrough transport headers, and never trust a convergence HEAD
// against such a target — a previous broken replica matches by ETag.
let Some(ssec_audit_required) = resolve_ssec_passthrough_gate(self.ssec, &tgt_client, &bucket, &object, &mut rinfo).await
else {
send_local_event(EventArgs {
event_name: EventName::ObjectReplicationNotTracked.to_string(),
bucket_name: bucket.clone(),
object: self.to_object_info(),
user_agent: "Internal: [Replication]".to_string(),
..Default::default()
});
return rinfo;
};
let versioned = ReplicationVersioningStore::prefix_enabled(&bucket, &object).await;
let version_suspended = ReplicationVersioningStore::prefix_suspended(&bucket, &object).await;
@@ -2985,18 +3094,20 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
}
let mut replication_action = replication_action;
match head_object_with_proxy_stats(
&bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
&object,
self.version_id.map(|v| v.to_string()),
)
.await
match head_object_for_worker(tgt_client.as_ref(), &tgt_client.bucket, &object, self.version_id.map(|v| v.to_string()))
.await
{
Ok(oi) => {
replication_action = replication_action_for_target_head(&object_info, &oi, self.op_type);
if replication_action == ReplicationAction::None {
// An SSE-C replica only counts as converged when the same
// HEAD proves its decryption material survived; a broken
// ciphertext copy from an earlier attempt matches by ETag.
if ssec_audit_required
&& !settle_ssec_passthrough_evidence(&oi, &tgt_client, &bucket, &object, &mut rinfo).await
{
return rinfo;
}
rinfo.replication_status = ReplicationStatusType::Completed;
rinfo.replication_resynced = true;
rinfo.replication_action = ReplicationAction::None;
@@ -3009,8 +3120,13 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
// Object not on target yet → fall through to PUT.
} else if is_version_id_format_mismatch(&e) {
// Version-ID format mismatch: retry without versionId and compare ETags.
match head_object_fallback(&bucket, &tgt_client, &object).await {
match head_object_fallback(&tgt_client, &object).await {
Ok(Some(oi)) if replication_etags_match(object_info.etag.as_deref(), oi.e_tag.as_deref()) => {
if ssec_audit_required
&& !settle_ssec_passthrough_evidence(&oi, &tgt_client, &bucket, &object, &mut rinfo).await
{
return rinfo;
}
rinfo.replication_status = ReplicationStatusType::Completed;
rinfo.replication_resynced = true;
rinfo.replication_action = ReplicationAction::None;
@@ -3085,7 +3201,6 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
}
};
let has_tagging_replication = !put_opts.user_tags.is_empty();
if let Some(err) = if is_multipart {
drop(gr);
let result = replicate_object_with_multipart(MultipartReplicationContext {
@@ -3100,10 +3215,6 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
put_opts,
})
.await;
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(&bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
} else {
gr.stream = wrap_with_bandwidth_monitor(gr.stream, &put_opts, &bucket, &rinfo.arn);
@@ -3119,10 +3230,6 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
)
})
.map_err(|e| std::io::Error::other(e.to_string()));
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(&bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
} {
rinfo.replication_status = ReplicationStatusType::Failed;
@@ -3144,6 +3251,15 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
return rinfo;
}
// First SSE-C passthrough PUT against this target: verify the replica
// kept its decryption material before reporting COMPLETED.
if ssec_audit_required
&& !audit_ssec_passthrough_replica(&tgt_client, &bucket, &object, self.version_id.map(|v| v.to_string()), &mut rinfo)
.await
{
return rinfo;
}
rinfo.replication_status = ReplicationStatusType::Completed;
rinfo
@@ -3166,6 +3282,21 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
return rinfo;
}
// N2 fail-closed: see the gate in `replicate_object` — the same policy
// applies to the metadata/existing-object transport.
let Some(ssec_audit_required) = resolve_ssec_passthrough_gate(self.ssec, &tgt_client, &bucket, &object, &mut rinfo).await
else {
send_local_event(EventArgs {
event_name: EventName::ObjectReplicationNotTracked.to_string(),
bucket_name: bucket.clone(),
object: self.to_object_info(),
user_agent: "Internal: [Replication]".to_string(),
..Default::default()
});
rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs();
return rinfo;
};
let versioned = ReplicationVersioningStore::prefix_enabled(&bucket, &object).await;
let version_suspended = ReplicationVersioningStore::prefix_suspended(&bucket, &object).await;
@@ -3204,8 +3335,19 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
let _sopts = replicate_all_stat_options(&object_info, &bucket, &tgt_client);
let Some((replication_action, object_info)) =
resolve_replicate_all_action(self, &tgt_client, &bucket, &object, object_info, start_time, &mut rinfo).await
let Some((replication_action, object_info)) = resolve_replicate_all_action(
ReplicateAllActionContext {
roi: self,
tgt_client: &tgt_client,
bucket: &bucket,
object: &object,
start_time,
ssec_audit_required,
},
object_info,
&mut rinfo,
)
.await
else {
return rinfo;
};
@@ -3260,6 +3402,16 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
return rinfo;
}
// First SSE-C passthrough PUT against this target: verify the replica
// kept its decryption material before reporting COMPLETED.
if ssec_audit_required
&& !audit_ssec_passthrough_replica(&tgt_client, &bucket, &object, self.version_id.map(|v| v.to_string()), &mut rinfo)
.await
{
rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs();
return rinfo;
}
rinfo
}
@@ -3472,33 +3624,49 @@ fn apply_replication_resync_timestamp(rinfo: &mut ReplicatedTargetInfo, reset_id
rinfo.replication_resynced = true;
}
/// Borrowed inputs for [`resolve_replicate_all_action`].
struct ReplicateAllActionContext<'a> {
roi: &'a ReplicateObjectInfo,
tgt_client: &'a Arc<TargetClient>,
bucket: &'a str,
object: &'a str,
start_time: OffsetDateTime,
/// N2: the target's SSE-C passthrough capability is still `Unknown`, so a
/// converged-looking replica must additionally prove its SSE-C material
/// survived before the comparison may settle COMPLETED.
ssec_audit_required: bool,
}
/// Compare the source object against the target via HEAD and decide which
/// replication action is still required. Returns `None` after fully settling
/// `rinfo` when replication must stop here — either because the target already
/// matches or because the comparison failed.
async fn resolve_replicate_all_action(
roi: &ReplicateObjectInfo,
tgt_client: &Arc<TargetClient>,
bucket: &str,
object: &str,
ctx: ReplicateAllActionContext<'_>,
object_info: ObjectInfo,
start_time: OffsetDateTime,
rinfo: &mut ReplicatedTargetInfo,
) -> Option<(ReplicationAction, ObjectInfo)> {
let replication_action;
match head_object_with_proxy_stats(
let ReplicateAllActionContext {
roi,
tgt_client,
bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
object,
roi.version_id.map(|v| v.to_string()),
)
.await
{
start_time,
ssec_audit_required,
} = ctx;
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) => {
replication_action = replication_action_for_target_head(&object_info, &oi, roi.op_type);
rinfo.replication_status = ReplicationStatusType::Completed;
if replication_action == ReplicationAction::None {
// An SSE-C replica only counts as converged when the same HEAD
// proves its decryption material survived; a broken ciphertext
// copy from an earlier attempt matches by ETag.
if ssec_audit_required && !settle_ssec_passthrough_evidence(&oi, tgt_client, bucket, object, rinfo).await {
rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs();
return None;
}
if roi.op_type == ReplicationType::ExistingObject
&& replication_target_head_is_newer_null_version(&object_info, &oi)
{
@@ -3545,9 +3713,15 @@ async fn resolve_replicate_all_action(
Err(e) => {
if is_version_id_format_mismatch(&e) {
// Version-ID format mismatch: retry without versionId and compare ETags.
match head_object_fallback(bucket, tgt_client, object).await {
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()) {
if ssec_audit_required
&& !settle_ssec_passthrough_evidence(&oi, tgt_client, bucket, object, rinfo).await
{
rinfo.duration = (OffsetDateTime::now_utc() - start_time).unsigned_abs();
return None;
}
ReplicationAction::None
} else {
ReplicationAction::All
@@ -3668,7 +3842,6 @@ async fn replicate_all_payload_to_target<S: ReplicationObjectIO>(
ctx: ReplicateAllPayloadContext<'_, S>,
mut gr: GetObjectReader,
) -> Option<std::io::Error> {
let has_tagging_replication = !ctx.put_opts.user_tags.is_empty();
if ctx.is_multipart {
drop(gr);
let result = replicate_object_with_multipart(MultipartReplicationContext {
@@ -3683,10 +3856,6 @@ async fn replicate_all_payload_to_target<S: ReplicationObjectIO>(
put_opts: ctx.put_opts,
})
.await;
record_proxy_request(ctx.bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(ctx.bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
} else {
gr.stream = wrap_with_bandwidth_monitor(gr.stream, &ctx.put_opts, ctx.bucket, ctx.arn);
@@ -3703,10 +3872,6 @@ async fn replicate_all_payload_to_target<S: ReplicationObjectIO>(
)
})
.map_err(|e| std::io::Error::other(e.to_string()));
record_proxy_request(ctx.bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(ctx.bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
}
}
@@ -1161,6 +1161,31 @@ mod tests {
assert!(all.contains_key("proxy-only-bucket"));
}
/// Pins the read-proxy metric contract (backlog#1675 P1-5): the API
/// strings the GET/HEAD/Tagging proxy paths record map onto the
/// get/head/tagging totals, and only unexpected failures raise the
/// failed counters.
#[tokio::test]
async fn test_proxy_stats_map_read_proxy_apis_to_totals() {
let stats = ReplicationStats::new();
stats.inc_proxy("proxy-bucket", "GetObject", false).await;
stats.inc_proxy("proxy-bucket", "GetObject", true).await;
stats.inc_proxy("proxy-bucket", "HeadObject", false).await;
stats.inc_proxy("proxy-bucket", "GetObjectTagging", false).await;
stats.inc_proxy("proxy-bucket", "PutObjectTagging", false).await;
stats.inc_proxy("proxy-bucket", "DeleteObjectTagging", true).await;
let metric = stats.get_proxy_stats("proxy-bucket").await;
assert_eq!(metric.get_total, 2);
assert_eq!(metric.get_failed, 1);
assert_eq!(metric.head_total, 1);
assert_eq!(metric.head_failed, 0);
assert_eq!(metric.get_tag_total, 1);
assert_eq!(metric.put_tag_total, 1);
assert_eq!(metric.delete_tag_total, 1);
assert_eq!(metric.delete_tag_failed, 1);
}
#[tokio::test]
async fn test_calculate_bucket_replication_stats_merges_resync_metrics() {
let stats = ReplicationStats::new();
@@ -36,7 +36,8 @@ use time::OffsetDateTime;
use time::format_description::well_known::Rfc3339;
pub(crate) use crate::bucket::bucket_target_sys::{
AdvancedPutOptions, PutObjectOptions, PutObjectPartOptions, RemoveObjectOptions, TargetClient,
AdvancedPutOptions, PutObjectOptions, PutObjectPartOptions, RemoveObjectOptions, SsecPassthroughCapability, TargetClient,
resolve_read_api_version_id,
};
#[cfg(test)]
pub(crate) use crate::bucket::target::BucketTarget;
@@ -65,6 +66,8 @@ static STANDARD_HEADERS: &[&str] = &[
];
const ERR_REPLICATION_ENCRYPTION_METADATA_UNSUPPORTED: &str = "replication source contains unsupported encryption metadata";
pub(crate) const ERR_REPLICATION_SSEC_PASSTHROUGH_UNSUPPORTED: &str = "replication target does not support SSE-C passthrough: the replica would lose its decryption material \
(run ?replication-check to re-probe)";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ReplicationSourceEncryption {
@@ -146,6 +149,54 @@ pub(crate) fn replication_object_is_ssec_encrypted(user_defined: &HashMap<String
rustfs_replication::is_ssec_encrypted(user_defined)
}
/// Fail-closed decision for an SSE-C passthrough replication attempt, derived
/// from the target's cached [`SsecPassthroughCapability`]. Pure so the policy
/// can migrate with the worker (M2) without dragging the cache along; the
/// caller computes `expired` from the cache record's age (see
/// `SSEC_PASSTHROUGH_CAPABILITY_TTL`).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SsecPassthroughGate {
/// Not an SSE-C object, or the target has a fresh proof that it preserves
/// the passthrough transport headers: replicate without a HEAD-back audit.
Proceed,
/// No usable verdict — first SSE-C attempt since the target was (re)built,
/// or the recorded verdict (in either direction) aged out: PUT, then HEAD
/// the replica back and require SSE-C evidence before reporting COMPLETED.
ProceedWithAudit,
/// The target was recently proven to drop the passthrough headers: do not
/// send the PUT, report FAILED (the object stays on the normal MRF retry
/// channel and re-audits once the verdict expires).
FailClosed,
}
pub(crate) fn ssec_passthrough_gate(ssec: bool, capability: SsecPassthroughCapability, expired: bool) -> SsecPassthroughGate {
if !ssec {
return SsecPassthroughGate::Proceed;
}
// An expired verdict — Supported or Unsupported — must be re-earned: a
// stale Unsupported would otherwise stick forever after a target upgrade,
// and a stale Supported would fail open after a backend swap behind the
// same endpoint.
if expired {
return SsecPassthroughGate::ProceedWithAudit;
}
match capability {
SsecPassthroughCapability::Supported => SsecPassthroughGate::Proceed,
SsecPassthroughCapability::Unknown => SsecPassthroughGate::ProceedWithAudit,
SsecPassthroughCapability::Unsupported => SsecPassthroughGate::FailClosed,
}
}
/// True when a replication-check HEAD of the replica proves the SSE-C
/// material survived passthrough: a RustFS target restores the transport
/// headers into the stored SSE-C keys and its HEAD echoes
/// `x-amz-server-side-encryption-customer-algorithm` (the replication-check
/// exemption skips key validation but not the metadata echo). A target that
/// dropped the headers stored a plain object and echoes nothing.
pub(crate) fn ssec_passthrough_evidence_present(head: &HeadObjectOutput) -> bool {
head.sse_customer_algorithm.as_deref().is_some_and(|algo| !algo.is_empty())
}
pub(crate) struct ReplicationTargetStore;
impl ReplicationTargetStore {
@@ -165,6 +216,17 @@ impl ReplicationTargetStore {
BucketTargetSys::get().mark_target_offline(target_client).await
}
/// Returns the cached verdict and whether it has outlived its TTL.
pub(crate) async fn ssec_passthrough_capability(arn: &str) -> (SsecPassthroughCapability, bool) {
BucketTargetSys::get().ssec_passthrough_capability(arn).await
}
pub(crate) async fn record_ssec_passthrough_capability(arn: &str, capability: SsecPassthroughCapability) {
BucketTargetSys::get()
.record_ssec_passthrough_capability(arn, capability)
.await
}
#[cfg(test)]
pub(crate) async fn register_test_target(target_client: &Arc<TargetClient>) {
BucketTargetSys::get().arn_remotes_map.write().await.insert(
@@ -898,6 +960,71 @@ mod tests {
}
}
/// N2 fail-closed policy: SSE-C replication may only proceed silently
/// against a target with a FRESH proof that it preserves the passthrough
/// transport headers. Unknown targets must be audited; freshly-flagged
/// dropping targets must never receive the PUT; an expired verdict in
/// EITHER direction must be re-earned through the audit — a sticky
/// Unsupported would outlive a target upgrade, and a sticky Supported
/// would fail open after a backend swap behind the same endpoint.
#[test]
fn ssec_passthrough_gate_is_fail_closed_and_ttl_bounded() {
for capability in [
SsecPassthroughCapability::Unknown,
SsecPassthroughCapability::Supported,
SsecPassthroughCapability::Unsupported,
] {
for expired in [false, true] {
assert_eq!(
ssec_passthrough_gate(false, capability, expired),
SsecPassthroughGate::Proceed,
"non-SSE-C objects must never be gated on the passthrough capability"
);
}
}
assert_eq!(
ssec_passthrough_gate(true, SsecPassthroughCapability::Supported, false),
SsecPassthroughGate::Proceed
);
assert_eq!(
ssec_passthrough_gate(true, SsecPassthroughCapability::Unknown, false),
SsecPassthroughGate::ProceedWithAudit
);
assert_eq!(
ssec_passthrough_gate(true, SsecPassthroughCapability::Unsupported, false),
SsecPassthroughGate::FailClosed
);
// Expiry flips both directions back to the audit.
assert_eq!(
ssec_passthrough_gate(true, SsecPassthroughCapability::Unsupported, true),
SsecPassthroughGate::ProceedWithAudit,
"an expired Unsupported verdict must allow a re-audit (upgraded target recovers without operator action)"
);
assert_eq!(
ssec_passthrough_gate(true, SsecPassthroughCapability::Supported, true),
SsecPassthroughGate::ProceedWithAudit,
"an expired Supported verdict must be re-proven (backend swap behind the same endpoint must not fail open)"
);
}
#[test]
fn ssec_passthrough_evidence_requires_customer_algorithm_echo() {
let with_evidence = HeadObjectOutput::builder().sse_customer_algorithm("AES256").build();
assert!(ssec_passthrough_evidence_present(&with_evidence));
let empty_algorithm = HeadObjectOutput::builder().sse_customer_algorithm("").build();
assert!(
!ssec_passthrough_evidence_present(&empty_algorithm),
"an empty echo is not evidence of preserved SSE-C material"
);
let without_evidence = HeadObjectOutput::builder().e_tag("\"abc\"").content_length(8).build();
assert!(
!ssec_passthrough_evidence_present(&without_evidence),
"a plain HEAD response must classify the target as having dropped the material"
);
}
#[test]
fn replication_put_options_adds_ssec_checksum_metadata() {
let metadata = HashMap::from([(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string())]);
+14
View File
@@ -277,6 +277,20 @@ 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,
/// True when the inbound request carried the
/// `{x-rustfs-,x-minio-}source-proxy-request` header family with the
/// value "true": the request was already proxied by a replication peer,
/// so this server must not proxy a local miss onward (anti-loop,
/// MinIO-compatible). The header only disables proxying — it grants no
/// capability — so no authorization gate is required to honor it.
pub proxy_request: bool,
/// True when the `source-proxy-request` header family was present at
/// all, regardless of value (MinIO's `ProxyHeaderSet`). A replication
/// peer sends `source-proxy-request: false` on its worker convergence
/// HEADs precisely so the receiver answers locally instead of proxying
/// back — otherwise a proxied 404->200 echo makes the worker believe the
/// object already converged and it never replicates it.
pub proxy_header_set: 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.